Merge commit 'refs/pr/530'

This commit is contained in:
fawney19
2026-05-28 12:53:45 +08:00
91 changed files with 4412 additions and 153 deletions
@@ -28,6 +28,12 @@ pub(crate) fn maybe_normalize_provider_private_sync_report_payload(
let mut normalized = payload.clone();
normalized.report_context = normalize_provider_private_report_context(Some(report_context));
if let (Some(body_json), Some(context)) = (
payload.body_json.as_ref(),
normalized.report_context.as_mut(),
) {
maybe_attach_gemini_cli_v1internal_credits_context(report_context, body_json, context);
}
if let Some(body_json) = payload.body_json.clone() {
normalized.body_json = normalize_provider_private_response_value(body_json, report_context);
@@ -55,6 +61,44 @@ pub(crate) fn maybe_normalize_provider_private_sync_report_payload(
Ok(Some(normalized))
}
fn maybe_attach_gemini_cli_v1internal_credits_context(
original_report_context: &Value,
body_json: &Value,
normalized_report_context: &mut Value,
) {
if !original_report_context
.get("envelope_name")
.and_then(Value::as_str)
.is_some_and(|value| value.eq_ignore_ascii_case("gemini_cli:v1internal"))
{
return;
}
let mut credits = serde_json::Map::new();
for (source, target) in [
("remainingCredits", "remainingCredits"),
("consumedCredits", "consumedCredits"),
("traceId", "traceId"),
] {
if let Some(value) = body_json
.get(source)
.cloned()
.filter(|value| !value.is_null())
{
credits.insert(target.to_string(), value);
}
}
if credits.is_empty() {
return;
}
if let Some(object) = normalized_report_context.as_object_mut() {
object.insert(
"gemini_cli_v1internal_credits".to_string(),
Value::Object(credits),
);
}
}
fn normalize_provider_private_stream_bytes(
report_context: &Value,
body: &[u8],
@@ -1859,6 +1859,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -198,6 +198,7 @@ mod tests {
}
}
})),
upstream_metadata: None,
decrypted_api_key: "sk-test".to_string(),
decrypted_auth_config: None,
},
@@ -254,6 +255,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "__placeholder__".to_string(),
decrypted_auth_config: None,
},
@@ -144,6 +144,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: String::new(),
decrypted_auth_config: None,
},
@@ -591,6 +591,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -0,0 +1,136 @@
use std::collections::BTreeMap;
use std::sync::Arc;
use serde_json::Value;
use crate::ai_serving::transport::gemini_cli::resolve_gemini_cli_project_id;
use crate::ai_serving::transport::{
build_gemini_cli_v1internal_request, build_standard_provider_request_headers,
GatewayProviderTransportSnapshot, GeminiCliRequestEnvelopeSupport,
StandardProviderRequestHeaders, StandardProviderRequestHeadersInput, GEMINI_CLI_USER_AGENT,
};
use crate::AppState;
pub(crate) enum GeminiCliV1InternalRequestError {
ProjectUnavailable,
EnvelopeUnsupported,
UpstreamUrlUnavailable,
HeaderRulesApplyFailed,
}
pub(crate) struct GeminiCliV1InternalRequestInput<'a> {
pub(crate) state: &'a AppState,
pub(crate) parts: &'a http::request::Parts,
pub(crate) transport: &'a Arc<GatewayProviderTransportSnapshot>,
pub(crate) trace_id: &'a str,
pub(crate) mapped_model: &'a str,
pub(crate) provider_api_format: &'a str,
pub(crate) auth_header: &'a str,
pub(crate) auth_value: &'a str,
pub(crate) request_headers: &'a http::HeaderMap,
pub(crate) original_request_body: &'a Value,
pub(crate) gemini_request_body: &'a Value,
pub(crate) upstream_is_stream: bool,
}
pub(crate) struct GeminiCliV1InternalRequest {
pub(crate) transport: Arc<GatewayProviderTransportSnapshot>,
pub(crate) body: Value,
pub(crate) headers: StandardProviderRequestHeaders,
pub(crate) upstream_url: String,
}
pub(crate) async fn build_gemini_cli_v1internal_provider_request(
input: GeminiCliV1InternalRequestInput<'_>,
) -> Result<GeminiCliV1InternalRequest, GeminiCliV1InternalRequestError> {
let payload = build_gemini_cli_v1internal_payload(
input.state,
input.transport,
input.trace_id,
input.mapped_model,
input.gemini_request_body,
)
.await?;
let upstream_url = crate::ai_serving::build_provider_transport_request_url_for_request_body(
&payload.transport,
input.provider_api_format,
Some(input.mapped_model),
input.upstream_is_stream,
input.parts.uri.query(),
None,
Some(&payload.body),
)
.ok_or(GeminiCliV1InternalRequestError::UpstreamUrlUnavailable)?;
let extra_headers =
BTreeMap::from([("user-agent".to_string(), GEMINI_CLI_USER_AGENT.to_string())]);
let headers = build_standard_provider_request_headers(StandardProviderRequestHeadersInput {
transport: &payload.transport,
provider_api_format: input.provider_api_format,
same_format: false,
headers: input.request_headers,
auth_header: input.auth_header,
auth_value: input.auth_value,
extra_headers: &extra_headers,
header_rules: payload.transport.endpoint.header_rules.as_ref(),
provider_request_body: &payload.body,
original_request_body: input.original_request_body,
upstream_is_stream: input.upstream_is_stream,
})
.ok_or(GeminiCliV1InternalRequestError::HeaderRulesApplyFailed)?;
Ok(GeminiCliV1InternalRequest {
transport: payload.transport,
body: payload.body,
headers,
upstream_url,
})
}
struct GeminiCliV1InternalPayload {
transport: Arc<GatewayProviderTransportSnapshot>,
body: Value,
}
async fn build_gemini_cli_v1internal_payload(
state: &AppState,
transport: &Arc<GatewayProviderTransportSnapshot>,
trace_id: &str,
mapped_model: &str,
gemini_request_body: &Value,
) -> Result<GeminiCliV1InternalPayload, GeminiCliV1InternalRequestError> {
let mut resolved_transport = Arc::clone(transport);
let project_id = match resolve_gemini_cli_project_id(&resolved_transport) {
Some(project_id) => Some(project_id),
None => match state
.hydrate_gemini_cli_project_metadata_for_transport(&resolved_transport)
.await
{
Some(hydrated) => {
let project_id = resolve_gemini_cli_project_id(&hydrated);
resolved_transport = Arc::new(hydrated);
project_id
}
None => None,
},
}
.ok_or(GeminiCliV1InternalRequestError::ProjectUnavailable)?;
let body = match build_gemini_cli_v1internal_request(
&project_id,
trace_id,
mapped_model,
gemini_request_body,
) {
GeminiCliRequestEnvelopeSupport::Supported(envelope) => envelope,
GeminiCliRequestEnvelopeSupport::Unsupported(_) => {
return Err(GeminiCliV1InternalRequestError::EnvelopeUnsupported);
}
};
Ok(GeminiCliV1InternalPayload {
transport: resolved_transport,
body,
})
}
@@ -12,6 +12,7 @@ mod candidate_transport_ranking_facts;
mod common;
mod decision;
mod decision_input;
mod gemini_cli;
mod materialization_policy;
mod passthrough;
mod plan_builders;
@@ -101,6 +101,11 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
super::super::ANTIGRAVITY_ENVELOPE_NAME,
parts.uri.path(),
);
} else if resolved.is_gemini_cli {
extra_fields.insert(
"envelope_name".to_string(),
json!(crate::ai_serving::transport::GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME),
);
}
let provider_api_format = resolved.provider_api_format.clone();
let effective_headers = input.effective_headers(&parts.headers);
@@ -144,7 +149,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
.and_then(serde_json::Value::as_bool)
.unwrap_or(false),
upstream_is_stream: resolved.upstream_is_stream,
has_envelope: resolved.is_kiro || resolved.is_antigravity,
has_envelope: resolved.is_kiro || resolved.is_antigravity || resolved.is_gemini_cli,
needs_conversion: false,
extra_fields,
}),
@@ -158,6 +163,7 @@ pub(crate) async fn maybe_build_local_same_format_provider_decision_payload_for_
let super::request::LocalSameFormatProviderCandidatePayloadParts {
transport,
is_antigravity: _,
is_gemini_cli: _,
is_kiro: _,
auth_header,
auth_value,
@@ -15,9 +15,11 @@ use crate::ai_serving::transport::antigravity::{
classify_local_antigravity_request_support, AntigravityEnvelopeRequestType,
AntigravityRequestEnvelopeSupport, AntigravityRequestSideSupport,
};
use crate::ai_serving::transport::gemini_cli::resolve_gemini_cli_project_id;
use crate::ai_serving::transport::{
build_grok_browser_headers, build_grok_upstream_url, build_same_format_provider_headers,
GrokHeaderInput, SameFormatProviderHeadersInput, GROK_CHAT_PATH,
build_gemini_cli_v1internal_request, build_grok_browser_headers, build_grok_upstream_url,
build_same_format_provider_headers, GeminiCliRequestEnvelopeSupport, GrokHeaderInput,
SameFormatProviderHeadersInput, GEMINI_CLI_USER_AGENT, GROK_CHAT_PATH,
};
use crate::ai_serving::{CandidateFailureDiagnostic, GatewayProviderTransportSnapshot};
use crate::{AppState, GatewayError};
@@ -91,6 +93,7 @@ pub(crate) fn resolve_same_format_provider_transport_unsupported_reason_for_trac
pub(crate) struct LocalSameFormatProviderCandidatePayloadParts {
pub(super) transport: Arc<GatewayProviderTransportSnapshot>,
pub(super) is_antigravity: bool,
pub(super) is_gemini_cli: bool,
pub(super) is_kiro: bool,
pub(super) auth_header: Option<String>,
pub(super) auth_value: Option<String>,
@@ -146,6 +149,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
)
.await?;
let body_json = redaction.body_json.as_ref();
let mut transport = Arc::clone(&prepared.transport);
let Some(mut base_provider_request_body) =
super::super::request::build_same_format_provider_request_body(
@@ -212,7 +216,7 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
let antigravity_auth = if prepared.is_antigravity {
match classify_local_antigravity_request_support(
&prepared.transport,
&transport,
&base_provider_request_body,
AntigravityEnvelopeRequestType::Agent,
) {
@@ -234,6 +238,38 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
} else {
None
};
let gemini_cli_project_id = if prepared.behavior.is_gemini_cli {
match resolve_gemini_cli_project_id(&transport) {
Some(project_id) => Some(project_id),
None => {
match state
.hydrate_gemini_cli_project_metadata_for_transport(&transport)
.await
{
Some(hydrated) => {
let project_id = resolve_gemini_cli_project_id(&hydrated);
transport = Arc::new(hydrated);
project_id
}
None => {
mark_skipped_local_same_format_provider_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"transport_auth_unavailable",
)
.await;
return Ok(None);
}
}
}
}
} else {
None
};
let provider_request_body = if let Some(antigravity_auth) = antigravity_auth.as_ref() {
match build_antigravity_safe_v1internal_request(
antigravity_auth,
@@ -263,6 +299,34 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
return Ok(None);
}
}
} else if let Some(project_id) = gemini_cli_project_id.as_deref() {
match build_gemini_cli_v1internal_request(
project_id,
trace_id,
&prepared.mapped_model,
&base_provider_request_body,
) {
GeminiCliRequestEnvelopeSupport::Supported(envelope) => envelope,
GeminiCliRequestEnvelopeSupport::Unsupported(_) => {
mark_skipped_local_same_format_provider_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_missing",
same_format_provider_request_body_failure_extra_data(
body_json,
attempt.eligible.provider_api_format.as_str(),
prepared.transport.endpoint.body_rules.as_ref(),
"gemini_cli_v1internal_envelope",
),
)
.await;
return Ok(None);
}
}
} else {
base_provider_request_body
};
@@ -273,14 +337,13 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
.provider_type
.trim()
.eq_ignore_ascii_case("grok");
let transport_profile =
crate::ai_serving::transport::resolve_transport_profile(&prepared.transport);
let transport_profile = crate::ai_serving::transport::resolve_transport_profile(&transport);
let upstream_url = if is_grok {
Some(build_grok_upstream_url(&prepared.transport, GROK_CHAT_PATH))
Some(build_grok_upstream_url(&transport, GROK_CHAT_PATH))
} else {
super::super::request::build_same_format_upstream_url(
parts,
&prepared.transport,
&transport,
&prepared.mapped_model,
prepared.provider_api_format.as_str(),
spec,
@@ -308,18 +371,21 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
return Ok(None);
};
let extra_headers = antigravity_auth
let mut extra_headers = antigravity_auth
.as_ref()
.map(build_antigravity_static_identity_headers)
.unwrap_or_default();
if prepared.behavior.is_gemini_cli {
extra_headers.insert("user-agent".to_string(), GEMINI_CLI_USER_AGENT.to_string());
}
let Some(mut provider_request_headers) = (if is_grok {
build_grok_browser_headers(GrokHeaderInput {
transport: &prepared.transport,
transport: &transport,
transport_profile: transport_profile.as_ref(),
request_headers: Some(effective_headers),
content_type: "application/json",
accept: "text/event-stream",
header_rules: prepared.transport.endpoint.header_rules.as_ref(),
header_rules: transport.endpoint.header_rules.as_ref(),
provider_request_body: &provider_request_body,
original_request_body: body_json,
})
@@ -328,12 +394,12 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
headers: effective_headers,
provider_request_body: &provider_request_body,
original_request_body: body_json,
header_rules: prepared.transport.endpoint.header_rules.as_ref(),
header_rules: transport.endpoint.header_rules.as_ref(),
behavior: prepared.behavior,
auth_header: prepared.auth_header.as_deref(),
auth_value: prepared.auth_value.as_deref(),
extra_headers: &extra_headers,
key_fingerprint: prepared.transport.key.fingerprint.as_ref(),
key_fingerprint: transport.key.fingerprint.as_ref(),
kiro_auth_config: prepared.kiro_auth.as_ref().map(|auth| &auth.auth_config),
kiro_machine_id: prepared
.kiro_auth
@@ -364,8 +430,9 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
);
Ok(Some(LocalSameFormatProviderCandidatePayloadParts {
transport: prepared.transport,
transport,
is_antigravity: prepared.is_antigravity,
is_gemini_cli: prepared.behavior.is_gemini_cli,
is_kiro: prepared.is_kiro,
auth_header: prepared.auth_header,
auth_value: prepared.auth_value,
@@ -439,6 +439,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "sk-upstream".to_string(),
decrypted_auth_config: None,
},
@@ -494,6 +495,46 @@ mod tests {
}
}
fn sample_gemini_cli_attempt(candidate_index: u32) -> LocalExecutionCandidateAttempt {
let mut transport = sample_transport("gemini:generate_content", "endpoint-gemini-cli");
transport.provider.provider_type = "gemini_cli".to_string();
transport.provider.name = "gemini".to_string();
transport.endpoint.base_url = "https://cloudcode-pa.googleapis.com".to_string();
transport.endpoint.custom_path = Some("/v1internal:{action}".to_string());
transport.endpoint.endpoint_kind = Some("generate_content".to_string());
transport.key.auth_type = "bearer".to_string();
transport.key.api_formats = Some(vec!["gemini:generate_content".to_string()]);
transport.key.global_priority_by_format = Some(json!({
"gemini:generate_content": 1,
}));
transport.key.upstream_metadata = Some(json!({
"gemini_cli": {
"project_id": "test-project"
}
}));
let mut candidate = sample_candidate("gemini:generate_content", "endpoint-gemini-cli");
candidate.provider_name = "gemini".to_string();
candidate.provider_type = "gemini_cli".to_string();
candidate.key_auth_type = "bearer".to_string();
candidate.selected_provider_model_name = "gemini-2.5-pro".to_string();
candidate.global_model_name = "gemini-2.5-pro".to_string();
LocalExecutionCandidateAttempt {
eligible: EligibleLocalExecutionCandidate {
kind: LocalExecutionCandidateKind::SingleKey,
candidate,
transport: Arc::new(transport),
provider_api_format: "gemini:generate_content".to_string(),
orchestration: LocalExecutionCandidateMetadata::default(),
ranking: None,
},
candidate_index,
retry_index: 0,
candidate_id: format!("candidate-{candidate_index}"),
}
}
fn claude_stream_spec() -> LocalStandardSpec {
LocalStandardSpec {
api_format: "claude:messages",
@@ -585,4 +626,90 @@ mod tests {
Some("bidirectional")
);
}
#[tokio::test]
async fn standard_family_wraps_gemini_cli_cross_format_body_in_v1internal_envelope() {
let state = crate::AppState::new().expect("state should build");
let request = http::Request::builder()
.method("POST")
.uri("/v1/chat/completions")
.header(http::header::CONTENT_TYPE, "application/json")
.body(())
.expect("request should build");
let (parts, _) = request.into_parts();
let body_json = json!({
"model": "gemini-2.5-pro",
"messages": [{"role": "user", "content": "hello"}],
"temperature": 0.2,
"stream": true
});
let mut input = sample_input();
input.requested_model = "gemini-2.5-pro".to_string();
let spec = LocalStandardSpec {
api_format: "openai:chat",
decision_kind: "openai_chat_stream",
report_kind: "openai_chat_stream_success",
family: LocalStandardSourceFamily::Standard,
mode: LocalStandardSourceMode::Chat,
require_streaming: true,
};
let payload = maybe_build_local_standard_decision_payload_for_candidate(
&state,
&parts,
"trace-gemini-cli-cross-format",
&body_json,
&input,
sample_gemini_cli_attempt(0),
spec,
)
.await
.expect("cross-format candidate should not fail routing mutation")
.expect("gemini_cli candidate should build a payload");
assert_eq!(
payload.upstream_url.as_deref(),
Some("https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse")
);
assert_eq!(
payload.execution_strategy.as_deref(),
Some("local_cross_format")
);
assert_eq!(
payload.provider_api_format.as_deref(),
Some("gemini:generate_content")
);
assert_eq!(
payload
.provider_request_headers
.get("user-agent")
.map(String::as_str),
Some(crate::ai_serving::transport::GEMINI_CLI_USER_AGENT)
);
let provider_body = payload
.provider_request_body
.as_ref()
.expect("provider request body should be present");
assert_eq!(provider_body["model"], "gemini-2.5-pro");
assert_eq!(provider_body["project"], "test-project");
assert_eq!(
provider_body["user_prompt_id"],
"trace-gemini-cli-cross-format"
);
assert!(provider_body.get("contents").is_none());
assert!(provider_body.get("generationConfig").is_none());
assert!(provider_body["request"].get("contents").is_some());
let report_context = payload
.report_context
.as_ref()
.expect("report context should be present");
assert_eq!(
report_context
.get("envelope_name")
.and_then(|value| value.as_str()),
Some(crate::ai_serving::transport::GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME)
);
}
}
@@ -12,6 +12,10 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, resolve_upstream_is_stream_for_provider,
};
use crate::ai_serving::planner::gemini_cli::{
build_gemini_cli_v1internal_provider_request, GeminiCliV1InternalRequestError,
GeminiCliV1InternalRequestInput,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
@@ -30,11 +34,12 @@ use crate::ai_serving::transport::{
build_openai_image_headers, build_openai_image_upstream_url,
build_standard_provider_request_headers, build_windsurf_cascade_headers,
build_windsurf_cascade_request_body, build_windsurf_cascade_upstream_url,
is_windsurf_provider_transport,
is_gemini_cli_provider_transport, is_windsurf_provider_transport,
local_windsurf_request_transport_unsupported_reason_with_network,
openai_image_transport_unsupported_reason, resolve_grok_session_auth,
resolve_openai_image_auth, GrokHeaderInput, ProviderOpenAiImageHeadersInput,
StandardProviderRequestHeadersInput, GROK_CHAT_PATH, WINDSURF_ENVELOPE_NAME,
StandardProviderRequestHeadersInput, GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME, GROK_CHAT_PATH,
WINDSURF_ENVELOPE_NAME,
};
use crate::ai_serving::{
build_openai_image_request_body_from_gemini_image_request, gemini_request_is_image_generation,
@@ -742,6 +747,31 @@ pub(crate) async fn resolve_local_standard_candidate_payload_parts(
.await);
}
let normalized_provider_api_format =
crate::ai_serving::normalize_api_format_alias(provider_api_format);
if normalized_provider_api_format == "gemini:generate_content"
&& is_gemini_cli_provider_transport(transport)
{
return Ok(build_gemini_cli_cross_format_payload_parts(
state,
parts,
trace_id,
body_json,
input,
attempt,
transport,
spec_metadata.api_format,
provider_api_format,
prepared_candidate.mapped_model,
prepared_candidate.auth_header,
prepared_candidate.auth_value,
provider_request_body,
upstream_is_stream,
redaction.redacted,
)
.await);
}
let upstream_url = match crate::ai_serving::planner::standard::build_standard_upstream_url(
parts,
transport,
@@ -845,6 +875,144 @@ fn apply_transport_request_body_semantics(
)
}
#[allow(clippy::too_many_arguments)]
async fn build_gemini_cli_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>,
client_api_format: &str,
provider_api_format: &str,
mapped_model: String,
auth_header: String,
auth_value: String,
gemini_request_body: Value,
upstream_is_stream: bool,
request_redacted: bool,
) -> Option<LocalStandardCandidatePayloadParts> {
let candidate = &attempt.eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
let resolved =
match build_gemini_cli_v1internal_provider_request(GeminiCliV1InternalRequestInput {
state,
parts,
transport,
trace_id,
mapped_model: &mapped_model,
provider_api_format,
auth_header: &auth_header,
auth_value: &auth_value,
request_headers: effective_headers,
original_request_body: original_body_json,
gemini_request_body: &gemini_request_body,
upstream_is_stream,
})
.await
{
Ok(resolved) => resolved,
Err(GeminiCliV1InternalRequestError::ProjectUnavailable) => {
mark_skipped_local_standard_candidate(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::EnvelopeUnsupported) => {
mark_skipped_local_standard_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
original_body_json,
client_api_format,
provider_api_format,
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::UpstreamUrlUnavailable) => {
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
client_api_format,
provider_api_format,
"standard_family_gemini_cli_url",
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::HeaderRulesApplyFailed) => {
mark_skipped_local_standard_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
attempt.candidate_index,
&attempt.candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
client_api_format,
provider_api_format,
"standard_family_gemini_cli_headers",
),
)
.await;
return None;
}
};
let mut provider_request_headers = resolved.headers.headers;
apply_codex_openai_responses_special_headers(
&mut provider_request_headers,
&resolved.body,
effective_headers,
resolved.transport.provider.provider_type.as_str(),
provider_api_format,
Some(trace_id),
resolved.transport.key.decrypted_auth_config.as_deref(),
);
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
Some(LocalStandardCandidatePayloadParts {
auth_header: resolved.headers.auth_header,
auth_value: resolved.headers.auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body: resolved.body,
provider_request_headers,
upstream_url: resolved.upstream_url,
upstream_is_stream,
envelope_name: Some(GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME),
transport: resolved.transport,
transport_profile: None,
request_redacted,
})
}
#[allow(clippy::too_many_arguments)]
async fn build_windsurf_cross_format_payload_parts(
state: &AppState,
@@ -68,6 +68,7 @@ fn sample_transport(base_url: &str, api_format: &str) -> GatewayProviderTranspor
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "__placeholder__".to_string(),
decrypted_auth_config: None,
},
@@ -13,6 +13,10 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, OPENAI_CHAT_STREAM_PLAN_KIND,
};
use crate::ai_serving::planner::gemini_cli::{
build_gemini_cli_v1internal_provider_request, GeminiCliV1InternalRequestError,
GeminiCliV1InternalRequestInput,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
@@ -38,9 +42,10 @@ use crate::ai_serving::transport::windsurf::{
use crate::ai_serving::transport::{
build_grok_browser_headers, build_grok_upstream_url, build_kiro_cross_format_upstream_url,
build_openai_image_headers, build_openai_image_upstream_url,
build_standard_provider_request_headers, openai_image_transport_unsupported_reason,
resolve_openai_image_auth, GrokHeaderInput, ProviderOpenAiImageHeadersInput,
StandardProviderRequestHeadersInput, GROK_CHAT_PATH,
build_standard_provider_request_headers, is_gemini_cli_provider_transport,
openai_image_transport_unsupported_reason, resolve_openai_image_auth, GrokHeaderInput,
ProviderOpenAiImageHeadersInput, StandardProviderRequestHeadersInput,
GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME, GROK_CHAT_PATH,
};
use crate::ai_serving::{
ai_local_execution_contract_for_formats, request_conversion_direct_auth,
@@ -652,6 +657,30 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
)
.await);
}
if provider_api_format == "gemini:generate_content"
&& is_gemini_cli_provider_transport(transport)
{
return Ok(build_gemini_cli_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,
redaction.redacted,
)
.await);
}
let Some(upstream_url) = build_cross_format_openai_chat_upstream_url(
parts,
@@ -752,6 +781,157 @@ pub(crate) async fn resolve_local_openai_chat_candidate_payload_parts(
}))
}
#[allow(clippy::too_many_arguments)]
async fn build_gemini_cli_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,
gemini_request_body: Value,
upstream_is_stream: bool,
request_redacted: bool,
) -> Option<LocalOpenAiChatCandidatePayloadParts> {
let candidate = &eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
let resolved =
match build_gemini_cli_v1internal_provider_request(GeminiCliV1InternalRequestInput {
state,
parts,
transport,
trace_id,
mapped_model: &mapped_model,
provider_api_format,
auth_header: &auth_header,
auth_value: &auth_value,
request_headers: effective_headers,
original_request_body: original_body_json,
gemini_request_body: &gemini_request_body,
upstream_is_stream,
})
.await
{
Ok(resolved) => resolved,
Err(GeminiCliV1InternalRequestError::ProjectUnavailable) => {
mark_skipped_local_openai_chat_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::EnvelopeUnsupported) => {
mark_skipped_local_openai_chat_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
original_body_json,
"openai:chat",
provider_api_format,
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::UpstreamUrlUnavailable) => {
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
"openai:chat",
provider_api_format,
"openai_chat_gemini_cli_url",
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::HeaderRulesApplyFailed) => {
mark_skipped_local_openai_chat_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
"openai:chat",
provider_api_format,
"openai_chat_gemini_cli_headers",
),
)
.await;
return None;
}
};
let mut provider_request_headers = resolved.headers.headers;
apply_codex_openai_responses_special_headers(
&mut provider_request_headers,
&resolved.body,
effective_headers,
resolved.transport.provider.provider_type.as_str(),
provider_api_format,
Some(trace_id),
resolved.transport.key.decrypted_auth_config.as_deref(),
);
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
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()
};
let (execution_strategy, conversion_mode) =
ai_local_execution_contract_for_formats("openai:chat", provider_api_format);
Some(LocalOpenAiChatCandidatePayloadParts {
client_api_format: "openai:chat".to_string(),
auth_header: resolved.headers.auth_header,
auth_value: resolved.headers.auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body: resolved.body,
provider_request_headers,
upstream_url: resolved.upstream_url,
execution_strategy,
conversion_mode,
report_kind: resolved_report_kind,
envelope_name: Some(GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME),
transport: resolved.transport,
request_redacted,
transport_profile: None,
image_request_summary: None,
})
}
#[allow(clippy::too_many_arguments)]
async fn resolve_openai_chat_to_openai_image_payload_parts(
state: &AppState,
@@ -1517,6 +1697,218 @@ async fn build_kiro_openai_chat_cross_format_payload_parts(
#[cfg(test)]
mod tests {
use super::*;
use aether_provider_transport::snapshot::{
GatewayProviderTransportEndpoint, GatewayProviderTransportKey,
GatewayProviderTransportProvider, GatewayProviderTransportSnapshot,
};
use aether_scheduler_core::SchedulerMinimalCandidateSelectionCandidate;
fn sample_auth_snapshot() -> crate::ai_serving::GatewayAuthApiKeySnapshot {
crate::ai_serving::GatewayAuthApiKeySnapshot {
user_id: "user-1".to_string(),
username: "alice".to_string(),
email: None,
user_role: "user".to_string(),
user_auth_source: "local".to_string(),
user_is_active: true,
user_is_deleted: false,
user_rate_limit: None,
user_allowed_providers: None,
user_allowed_api_formats: None,
user_allowed_models: None,
api_key_id: "api-key-1".to_string(),
api_key_name: Some("default".to_string()),
api_key_is_active: true,
api_key_is_locked: false,
api_key_is_standalone: false,
api_key_rate_limit: None,
api_key_concurrent_limit: None,
api_key_expires_at_unix_secs: None,
api_key_allowed_providers: None,
api_key_allowed_api_formats: None,
api_key_allowed_models: None,
api_key_ip_rules: None,
currently_usable: true,
}
}
fn sample_input() -> LocalOpenAiChatDecisionInput {
LocalOpenAiChatDecisionInput {
auth_context: crate::ai_serving::ExecutionRuntimeAuthContext {
user_id: "user-1".to_string(),
api_key_id: "api-key-1".to_string(),
username: Some("alice".to_string()),
api_key_name: Some("default".to_string()),
balance_remaining: Some(10.0),
access_allowed: true,
api_key_is_standalone: false,
},
requested_model: "gemini-2.5-pro".to_string(),
auth_snapshot: sample_auth_snapshot(),
required_capabilities: None,
request_auth_channel: None,
client_session_affinity: None,
routing_policy: None,
routing_trace_seed: None,
routing_context: None,
}
}
fn sample_gemini_cli_transport() -> GatewayProviderTransportSnapshot {
GatewayProviderTransportSnapshot {
provider: GatewayProviderTransportProvider {
id: "provider-1".to_string(),
name: "gemini".to_string(),
provider_type: "gemini_cli".to_string(),
website: None,
is_active: true,
keep_priority_on_conversion: false,
enable_format_conversion: true,
concurrent_limit: None,
max_retries: None,
proxy: None,
request_timeout_secs: None,
stream_first_byte_timeout_secs: None,
config: None,
},
endpoint: GatewayProviderTransportEndpoint {
id: "endpoint-1".to_string(),
provider_id: "provider-1".to_string(),
api_format: "gemini:generate_content".to_string(),
api_family: Some("gemini".to_string()),
endpoint_kind: Some("generate_content".to_string()),
is_active: true,
base_url: "https://cloudcode-pa.googleapis.com".to_string(),
header_rules: None,
body_rules: None,
max_retries: None,
custom_path: Some("/v1internal:{action}".to_string()),
config: None,
format_acceptance_config: None,
proxy: None,
},
key: GatewayProviderTransportKey {
id: "key-1".to_string(),
provider_id: "provider-1".to_string(),
name: "key".to_string(),
auth_type: "bearer".to_string(),
is_active: true,
api_formats: Some(vec!["gemini:generate_content".to_string()]),
auth_type_by_format: None,
allow_auth_channel_mismatch_formats: None,
allowed_models: None,
capabilities: None,
rate_multipliers: None,
global_priority_by_format: Some(json!({
"gemini:generate_content": 1,
})),
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: Some(json!({
"gemini_cli": {
"project_id": "test-project"
}
})),
decrypted_api_key: "oauth-access-token".to_string(),
decrypted_auth_config: None,
},
}
}
fn sample_gemini_cli_eligible() -> EligibleLocalExecutionCandidate {
EligibleLocalExecutionCandidate {
kind: crate::ai_serving::planner::candidate_resolution::LocalExecutionCandidateKind::SingleKey,
candidate: SchedulerMinimalCandidateSelectionCandidate {
provider_id: "provider-1".to_string(),
provider_name: "gemini".to_string(),
provider_type: "gemini_cli".to_string(),
provider_priority: 1,
endpoint_id: "endpoint-1".to_string(),
endpoint_api_format: "gemini:generate_content".to_string(),
key_id: "key-1".to_string(),
key_name: "key".to_string(),
key_auth_type: "bearer".to_string(),
key_internal_priority: 1,
key_global_priority_for_format: Some(1),
key_capabilities: None,
model_id: "model-1".to_string(),
global_model_id: "global-model-1".to_string(),
global_model_name: "gemini-2.5-pro".to_string(),
selected_provider_model_name: "gemini-2.5-pro".to_string(),
mapping_matched_model: None,
},
transport: Arc::new(sample_gemini_cli_transport()),
provider_api_format: "gemini:generate_content".to_string(),
orchestration: crate::orchestration::LocalExecutionCandidateMetadata::default(),
ranking: None,
}
}
#[tokio::test]
async fn openai_chat_to_gemini_cli_wraps_cross_format_body_in_v1internal_envelope() {
let state = AppState::new().expect("state should build");
let request = http::Request::builder()
.method("POST")
.uri("/v1/chat/completions")
.header(http::header::CONTENT_TYPE, "application/json")
.body(())
.expect("request should build");
let (parts, _) = request.into_parts();
let body_json = json!({
"model": "gemini-2.5-pro",
"messages": [{"role": "user", "content": "hello"}],
"generationConfig": {"temperature": 0.2},
"stream": true
});
let payload = resolve_local_openai_chat_candidate_payload_parts(
&state,
&parts,
"trace-openai-chat-gemini-cli",
&body_json,
&sample_input(),
&sample_gemini_cli_eligible(),
0,
"candidate-0",
OPENAI_CHAT_STREAM_PLAN_KIND,
"openai_chat_stream_success",
true,
)
.await
.expect("candidate resolution should not fail")
.expect("gemini_cli candidate should build a payload");
assert_eq!(
payload.upstream_url,
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
payload.envelope_name,
Some(GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME)
);
assert_eq!(
payload
.provider_request_headers
.get("user-agent")
.map(String::as_str),
Some(crate::ai_serving::transport::GEMINI_CLI_USER_AGENT)
);
assert_eq!(payload.provider_request_body["model"], "gemini-2.5-pro");
assert_eq!(payload.provider_request_body["project"], "test-project");
assert_eq!(
payload.provider_request_body["user_prompt_id"],
"trace-openai-chat-gemini-cli"
);
assert!(payload.provider_request_body.get("contents").is_none());
assert!(payload
.provider_request_body
.get("generationConfig")
.is_none());
assert!(payload.provider_request_body["request"]
.get("contents")
.is_some());
}
#[test]
fn chatgpt_web_chat_image_bridge_body_uses_internal_web_shape() {
@@ -103,6 +103,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -14,6 +14,10 @@ use crate::ai_serving::planner::common::{
endpoint_config_forces_body_stream_field, enforce_provider_body_stream_policy,
request_requires_body_stream_field, resolve_upstream_is_stream_for_provider,
};
use crate::ai_serving::planner::gemini_cli::{
build_gemini_cli_v1internal_provider_request, GeminiCliV1InternalRequestError,
GeminiCliV1InternalRequestInput,
};
use crate::ai_serving::planner::redaction::{
request_identity_response_encoding_when_redacted, resolve_provider_chat_pii_redaction,
};
@@ -44,11 +48,12 @@ use crate::ai_serving::transport::{
build_openai_image_headers, build_openai_image_upstream_url,
build_standard_provider_request_headers, build_windsurf_cascade_headers,
build_windsurf_cascade_request_body, build_windsurf_cascade_upstream_url,
is_windsurf_provider_transport, local_standard_transport_unsupported_reason_with_network,
is_gemini_cli_provider_transport, is_windsurf_provider_transport,
local_standard_transport_unsupported_reason_with_network,
local_windsurf_request_transport_unsupported_reason_with_network,
openai_image_transport_unsupported_reason, resolve_openai_image_auth, GrokHeaderInput,
ProviderOpenAiImageHeadersInput, StandardProviderRequestHeadersInput, GROK_CHAT_PATH,
WINDSURF_ENVELOPE_NAME,
ProviderOpenAiImageHeadersInput, StandardProviderRequestHeadersInput,
GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME, GROK_CHAT_PATH, WINDSURF_ENVELOPE_NAME,
};
use crate::ai_serving::{
ai_local_execution_contract_for_formats, request_conversion_direct_auth,
@@ -498,6 +503,30 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
)
.await);
}
if provider_api_format == "gemini:generate_content"
&& is_gemini_cli_provider_transport(transport)
{
return Ok(build_gemini_cli_openai_responses_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,
redaction.redacted,
)
.await);
}
let Some(upstream_url) = (if is_grok && is_grok_text_provider_api_format(provider_api_format) {
Some(build_grok_upstream_url(transport, GROK_CHAT_PATH))
@@ -675,6 +704,152 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
}))
}
#[allow(clippy::too_many_arguments)]
async fn build_gemini_cli_openai_responses_payload_parts(
state: &AppState,
parts: &http::request::Parts,
trace_id: &str,
original_body_json: &serde_json::Value,
input: &LocalOpenAiResponsesDecisionInput,
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,
gemini_request_body: Value,
upstream_is_stream: bool,
request_redacted: bool,
) -> Option<LocalOpenAiResponsesCandidatePayloadParts> {
let candidate = &eligible.candidate;
let effective_headers = input.effective_headers(&parts.headers);
let resolved =
match build_gemini_cli_v1internal_provider_request(GeminiCliV1InternalRequestInput {
state,
parts,
transport,
trace_id,
mapped_model: &mapped_model,
provider_api_format,
auth_header: &auth_header,
auth_value: &auth_value,
request_headers: effective_headers,
original_request_body: original_body_json,
gemini_request_body: &gemini_request_body,
upstream_is_stream,
})
.await
{
Ok(resolved) => resolved,
Err(GeminiCliV1InternalRequestError::ProjectUnavailable) => {
mark_skipped_local_openai_responses_candidate(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_auth_unavailable",
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::EnvelopeUnsupported) => {
mark_skipped_local_openai_responses_candidate_with_extra_data(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"provider_request_body_build_failed",
request_body_build_failure_extra_data(
original_body_json,
client_api_format,
provider_api_format,
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::UpstreamUrlUnavailable) => {
mark_skipped_local_openai_responses_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"upstream_url_missing",
CandidateFailureDiagnostic::upstream_url_missing(
client_api_format,
provider_api_format,
"openai_responses_gemini_cli_url",
),
)
.await;
return None;
}
Err(GeminiCliV1InternalRequestError::HeaderRulesApplyFailed) => {
mark_skipped_local_openai_responses_candidate_with_failure_diagnostic(
state,
input,
trace_id,
candidate,
candidate_index,
candidate_id,
"transport_header_rules_apply_failed",
CandidateFailureDiagnostic::header_rules_apply_failed(
client_api_format,
provider_api_format,
"openai_responses_gemini_cli_headers",
),
)
.await;
return None;
}
};
let mut provider_request_headers = resolved.headers.headers;
apply_codex_openai_responses_special_headers(
&mut provider_request_headers,
&resolved.body,
effective_headers,
resolved.transport.provider.provider_type.as_str(),
provider_api_format,
Some(trace_id),
resolved.transport.key.decrypted_auth_config.as_deref(),
);
request_identity_response_encoding_when_redacted(
&mut provider_request_headers,
request_redacted,
);
let (execution_strategy, conversion_mode) =
ai_local_execution_contract_for_formats(client_api_format, provider_api_format);
Some(LocalOpenAiResponsesCandidatePayloadParts {
auth_header: resolved.headers.auth_header,
auth_value: resolved.headers.auth_value,
mapped_model,
provider_api_format: provider_api_format.to_string(),
provider_request_body: resolved.body,
provider_request_headers,
upstream_url: resolved.upstream_url,
execution_strategy,
conversion_mode,
is_antigravity: false,
envelope_name: Some(GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME),
upstream_is_stream,
transport: resolved.transport,
transport_profile: None,
image_request_summary: None,
request_redacted,
})
}
#[allow(clippy::too_many_arguments)]
async fn build_windsurf_openai_responses_payload_parts(
state: &AppState,
+25 -19
View File
@@ -18,6 +18,10 @@ pub(crate) mod grok {
pub(crate) use aether_provider_transport::grok::*;
}
pub(crate) mod gemini_cli {
pub(crate) use aether_provider_transport::gemini_cli::*;
}
pub(crate) mod oauth_refresh {
pub(crate) use aether_provider_transport::oauth_refresh::*;
}
@@ -62,14 +66,14 @@ pub(crate) use aether_provider_transport::{
apply_transport_request_body_semantics, body_rules_are_locally_supported,
body_rules_handle_path, body_rules_have_enabled_rules,
build_cross_format_openai_chat_upstream_url, build_cross_format_openai_responses_upstream_url,
build_gemini_files_headers, build_gemini_files_request_body, build_gemini_files_upstream_url,
build_grok_app_chat_body, build_grok_browser_headers, build_grok_upstream_url,
build_kiro_cross_format_upstream_url, build_local_openai_chat_upstream_url,
build_local_openai_responses_upstream_url, build_openai_image_headers,
build_openai_image_upstream_url, build_passthrough_headers, build_request_trace_proxy_value,
build_same_format_provider_headers, build_same_format_provider_request_body,
build_same_format_provider_upstream_url, build_standard_plan_fallback_headers,
build_standard_plan_fallback_openai_chat_url,
build_gemini_cli_v1internal_request, build_gemini_files_headers,
build_gemini_files_request_body, build_gemini_files_upstream_url, build_grok_app_chat_body,
build_grok_browser_headers, build_grok_upstream_url, build_kiro_cross_format_upstream_url,
build_local_openai_chat_upstream_url, build_local_openai_responses_upstream_url,
build_openai_image_headers, build_openai_image_upstream_url, build_passthrough_headers,
build_request_trace_proxy_value, build_same_format_provider_headers,
build_same_format_provider_request_body, build_same_format_provider_upstream_url,
build_standard_plan_fallback_headers, build_standard_plan_fallback_openai_chat_url,
build_standard_plan_fallback_openai_responses_url, build_standard_provider_request_headers,
build_transport_request_url, build_transport_request_url_for_request_body,
build_video_create_headers, build_video_create_request_body, build_video_create_upstream_url,
@@ -78,7 +82,8 @@ pub(crate) use aether_provider_transport::{
candidate_transport_pair_skip_reason, classify_same_format_provider_request_behavior,
ensure_upstream_auth_header, gemini_files_transport_unsupported_reason,
header_rules_are_locally_supported, header_rules_have_enabled_rules,
is_windsurf_provider_transport, local_gemini_transport_unsupported_reason_with_network,
is_gemini_cli_provider_transport, is_windsurf_provider_transport,
local_gemini_transport_unsupported_reason_with_network,
local_openai_chat_transport_unsupported_reason,
local_standard_transport_unsupported_reason_with_network,
local_windsurf_request_transport_unsupported_reason_with_network,
@@ -95,14 +100,15 @@ pub(crate) use aether_provider_transport::{
supports_local_generic_oauth_request_auth_resolution,
supports_local_oauth_request_auth_resolution, transport_proxy_is_locally_supported,
video_create_transport_unsupported_reason, CandidateTransportPolicyFacts,
GatewayProviderTransportSnapshot, GeminiFilesHeadersInput, GeminiFilesRequestBodyError,
GeminiFilesRequestBodyParts, GrokHeaderInput, LocalResolvedOAuthRequestAuth,
ProviderOpenAiImageHeadersInput, ProviderVideoCreateFamily, ProviderVideoCreateHeadersInput,
SameFormatProviderFamily, SameFormatProviderHeadersInput, SameFormatProviderRequestBehavior,
SameFormatProviderRequestBehaviorParams, SameFormatProviderRequestBodyInput,
SameFormatProviderUpstreamUrlParams, StandardPlanFallbackAcceptPolicy,
StandardPlanFallbackHeadersInput, StandardProviderRequestHeaders,
StandardProviderRequestHeadersInput, TransportRequestBodySemanticsError,
TransportRequestUrlParams, GROK_CHAT_PATH, GROK_INTERNAL_HEADER, GROK_RATE_LIMITS_PATH,
WINDSURF_ENVELOPE_NAME,
GatewayProviderTransportSnapshot, GeminiCliRequestEnvelopeSupport, GeminiFilesHeadersInput,
GeminiFilesRequestBodyError, GeminiFilesRequestBodyParts, GrokHeaderInput,
LocalResolvedOAuthRequestAuth, ProviderOpenAiImageHeadersInput, ProviderVideoCreateFamily,
ProviderVideoCreateHeadersInput, SameFormatProviderFamily, SameFormatProviderHeadersInput,
SameFormatProviderRequestBehavior, SameFormatProviderRequestBehaviorParams,
SameFormatProviderRequestBodyInput, SameFormatProviderUpstreamUrlParams,
StandardPlanFallbackAcceptPolicy, StandardPlanFallbackHeadersInput,
StandardProviderRequestHeaders, StandardProviderRequestHeadersInput,
TransportRequestBodySemanticsError, TransportRequestUrlParams, GEMINI_CLI_USER_AGENT,
GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME, GROK_CHAT_PATH, GROK_INTERNAL_HEADER,
GROK_RATE_LIMITS_PATH, WINDSURF_ENVELOPE_NAME,
};
+16 -5
View File
@@ -97,7 +97,7 @@ pub(super) fn classify_ai_public_route(
"openai:video",
true,
))
} else if is_gemini_models_route(normalized_path) {
} else if method == http::Method::POST && is_gemini_models_route(normalized_path) {
if normalized_path.ends_with(":predictLongRunning") {
Some(classified(
"ai_public",
@@ -136,7 +136,9 @@ pub(super) fn classify_ai_public_route(
true,
))
}
} else if is_gemini_operation_route(normalized_path) {
} else if is_gemini_operation_method(method, normalized_path)
&& is_gemini_operation_route(normalized_path)
{
Some(classified(
"ai_public",
"gemini",
@@ -144,9 +146,7 @@ pub(super) fn classify_ai_public_route(
"gemini:video",
true,
))
} else if (method == http::Method::POST && normalized_path == "/upload/v1beta/files")
|| normalized_path.starts_with("/v1beta/files")
{
} else if is_gemini_files_method(method, normalized_path) {
Some(classified(
"ai_public",
"gemini",
@@ -159,6 +159,17 @@ pub(super) fn classify_ai_public_route(
}
}
fn is_gemini_operation_method(method: &http::Method, normalized_path: &str) -> bool {
method == http::Method::GET
|| (method == http::Method::POST && normalized_path.ends_with(":cancel"))
}
fn is_gemini_files_method(method: &http::Method, normalized_path: &str) -> bool {
(method == http::Method::POST && normalized_path == "/upload/v1beta/files")
|| ((method == http::Method::GET || method == http::Method::DELETE)
&& normalized_path.starts_with("/v1beta/files"))
}
fn classify_antigravity_v1internal_route(
method: &http::Method,
normalized_path: &str,
@@ -219,6 +219,23 @@ fn classifies_gemini_generate_content_api_key_without_cli_marker() {
assert!(decision.is_execution_runtime_candidate());
}
#[test]
fn does_not_classify_gemini_options_preflight_as_ai_execution_route() {
let headers = headers(&[
("origin", "http://localhost:3000"),
("access-control-request-method", "POST"),
(
"access-control-request-headers",
"authorization,content-type",
),
]);
let uri: Uri = "/v1beta/models/gemini-3-flash-preview:streamGenerateContent?alt=sse"
.parse()
.expect("uri should parse");
assert!(classify_control_route(&http::Method::OPTIONS, &uri, &headers).is_none());
}
#[test]
fn classifies_gemini_embed_content_as_embedding_route() {
let headers = headers(&[("x-goog-api-key", "gemini-key")]);
@@ -4139,6 +4139,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -4252,6 +4253,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
+1
View File
@@ -196,6 +196,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -380,6 +380,17 @@ fn invalid_gemini_provider_success_message(
{
return None;
}
let normalized_body_json = report_context
.filter(|context| {
context
.get("has_envelope")
.and_then(Value::as_bool)
.unwrap_or(false)
})
.and_then(|context| {
crate::ai_serving::normalize_provider_private_response_value(body_json.clone(), context)
});
let body_json = normalized_body_json.as_ref().unwrap_or(body_json);
if crate::ai_serving::gemini_generate_content_response_has_visible_output(body_json) {
return None;
}
@@ -2520,6 +2531,39 @@ mod tests {
assert!(message.contains("visible model output"));
}
#[test]
fn invalid_gemini_provider_success_unwraps_gemini_cli_v1internal_envelope() {
let plan = test_gemini_chat_plan();
let report_context = json!({
"has_envelope": true,
"envelope_name": "gemini_cli:v1internal",
"provider_api_format": "gemini:generate_content",
});
let body = json!({
"response": {
"candidates": [{
"content": {
"role": "model",
"parts": [{"text": "Hello from Gemini CLI"}]
},
"finishReason": "STOP"
}]
},
"remainingCredits": 41,
"consumedCredits": 1,
"traceId": "trace-upstream-sync-1"
});
let message = invalid_gemini_provider_success_message(
&plan,
Some(&report_context),
StatusCode::OK.as_u16(),
Some(&body),
);
assert!(message.is_none());
}
#[tokio::test]
async fn sync_attempt_terminal_guard_marks_dropped_pending_attempt_cancelled() {
let usage_repository = Arc::new(InMemoryUsageReadRepository::default());
@@ -31,6 +31,7 @@ pub(super) struct AdminProviderOAuthBatchImportEntry {
pub user_id: Option<String>,
pub email: Option<String>,
pub account_name: Option<String>,
pub project_id: Option<String>,
pub sso_rw_token: Option<String>,
pub cf_cookies: Option<String>,
pub cf_clearance: Option<String>,
@@ -76,6 +77,20 @@ fn coerce_admin_provider_oauth_import_str(value: Option<&serde_json::Value>) ->
.map(ToOwned::to_owned)
}
fn coerce_admin_provider_oauth_import_project_id(
value: Option<&serde_json::Value>,
) -> Option<String> {
match value {
Some(serde_json::Value::Object(object)) => coerce_admin_provider_oauth_import_str(
object
.get("id")
.or_else(|| object.get("project_id"))
.or_else(|| object.get("projectId")),
),
other => coerce_admin_provider_oauth_import_str(other),
}
}
fn json_import_expiry_value(value: Option<&serde_json::Value>) -> Option<u64> {
let value = value?;
json_u64_value(Some(value)).or_else(|| {
@@ -175,6 +190,7 @@ fn extract_admin_provider_oauth_batch_import_entry(
user_id: grok_cookie_value(raw_token, "x-userid"),
email: None,
account_name: None,
project_id: None,
sso_rw_token: grok_cookie_value(raw_token, "sso-rw"),
cf_cookies: grok_cookie_profile(raw_token),
cf_clearance: grok_cookie_value(raw_token, "cf_clearance"),
@@ -308,6 +324,13 @@ fn extract_admin_provider_oauth_batch_import_entry(
.get("account_name")
.or_else(|| object.get("accountName")),
);
let project_id = coerce_admin_provider_oauth_import_project_id(
object
.get("project_id")
.or_else(|| object.get("projectId"))
.or_else(|| object.get("cloudaicompanionProject"))
.or_else(|| object.get("cloudAiCompanionProject")),
);
let sso_rw_token = coerce_admin_provider_oauth_import_str(
object
.get("sso_rw_token")
@@ -355,6 +378,7 @@ fn extract_admin_provider_oauth_batch_import_entry(
user_id,
email,
account_name,
project_id,
sso_rw_token,
cf_cookies,
cf_clearance,
@@ -445,6 +469,7 @@ fn parse_error_entry(error: String) -> AdminProviderOAuthBatchImportEntry {
user_id: None,
email: None,
account_name: None,
project_id: None,
sso_rw_token: None,
cf_cookies: None,
cf_clearance: None,
@@ -464,6 +489,19 @@ pub(super) fn apply_admin_provider_oauth_batch_import_hints(
auth_config: &mut serde_json::Map<String, serde_json::Value>,
) {
let provider_type = provider_type.trim().to_ascii_lowercase();
if provider_type == "gemini_cli" {
if let Some(project_id) = entry.project_id.as_ref() {
auth_config
.entry("project_id".to_string())
.or_insert_with(|| json!(project_id));
}
if let Some(plan_type) = entry.plan_type.as_ref() {
auth_config
.entry("plan_type".to_string())
.or_insert_with(|| json!(plan_type));
}
return;
}
if !matches!(provider_type.as_str(), "codex" | "chatgpt_web" | "grok") {
return;
}
@@ -614,7 +652,10 @@ pub(super) fn build_admin_provider_oauth_batch_task_state(
#[cfg(test)]
mod tests {
use super::parse_admin_provider_oauth_batch_import_entries;
use super::{
apply_admin_provider_oauth_batch_import_hints,
parse_admin_provider_oauth_batch_import_entries,
};
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use serde_json::json;
@@ -750,6 +791,38 @@ mod tests {
assert_eq!(entries[0].pool_tier.as_deref(), Some("heavy"));
}
#[test]
fn parses_gemini_cli_project_id_hint() {
let entries = parse_admin_provider_oauth_batch_import_entries(
"gemini_cli",
r#"[{"refresh_token":"rt-1","projectId":"project-gemini-cli-1","planType":"free"}]"#,
);
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].refresh_token.as_deref(), Some("rt-1"));
assert_eq!(
entries[0].project_id.as_deref(),
Some("project-gemini-cli-1")
);
assert_eq!(entries[0].plan_type.as_deref(), Some("free"));
}
#[test]
fn applies_gemini_cli_project_id_hint_to_auth_config() {
let entries = parse_admin_provider_oauth_batch_import_entries(
"gemini_cli",
r#"{"refreshToken":"rt-1","cloudaicompanionProject":{"id":"project-gemini-cli-2"}}"#,
);
let mut auth_config = serde_json::Map::new();
apply_admin_provider_oauth_batch_import_hints("gemini_cli", &entries[0], &mut auth_config);
assert_eq!(
auth_config.get("project_id"),
Some(&json!("project-gemini-cli-2"))
);
}
#[test]
fn parses_windsurf_json_credentials_for_native_import() {
let entries = parse_admin_provider_oauth_batch_import_entries(
@@ -4,6 +4,7 @@ use std::pin::Pin;
use super::antigravity::refresh_antigravity_provider_quota_locally;
use super::chatgpt_web::refresh_chatgpt_web_provider_quota_locally;
use super::codex::refresh_codex_provider_quota_locally;
use super::gemini_cli::refresh_gemini_cli_provider_quota_locally;
use super::grok::refresh_grok_provider_quota_locally;
use super::kiro::refresh_kiro_provider_quota_locally;
use super::windsurf::refresh_windsurf_provider_quota_locally;
@@ -35,6 +36,10 @@ const PROVIDER_QUOTA_REFRESH_HANDLERS: &[(&str, ProviderQuotaRefreshHandler)] =
refresh_chatgpt_web_provider_quota_locally_boxed,
),
("codex", refresh_codex_provider_quota_locally_boxed),
(
"gemini_cli",
refresh_gemini_cli_provider_quota_locally_boxed,
),
("grok", refresh_grok_provider_quota_locally_boxed),
("kiro", refresh_kiro_provider_quota_locally_boxed),
("windsurf", refresh_windsurf_provider_quota_locally_boxed),
@@ -106,6 +111,22 @@ fn refresh_codex_provider_quota_locally_boxed<'a>(
))
}
fn refresh_gemini_cli_provider_quota_locally_boxed<'a>(
state: &'a AdminAppState<'a>,
provider: &'a StoredProviderCatalogProvider,
endpoint: &'a StoredProviderCatalogEndpoint,
keys: Vec<StoredProviderCatalogKey>,
proxy_override: Option<ProxySnapshot>,
) -> ProviderQuotaRefreshFuture<'a> {
Box::pin(refresh_gemini_cli_provider_quota_locally(
state,
provider,
endpoint,
keys,
proxy_override,
))
}
fn refresh_kiro_provider_quota_locally_boxed<'a>(
state: &'a AdminAppState<'a>,
provider: &'a StoredProviderCatalogProvider,
@@ -0,0 +1,271 @@
use super::shared::{
build_provider_quota_execution_plan, build_quota_snapshot_payload,
default_provider_quota_execution_timeouts, execute_provider_quota_plan,
extract_execution_error_message, oauth_refresh_auto_removed_result,
persist_provider_quota_refresh_state, quota_key_auto_removed,
quota_refresh_success_invalid_state, ProviderQuotaExecutionOutcome,
};
use crate::handlers::admin::request::{AdminAppState, AdminGatewayProviderTransportSnapshot};
use crate::GatewayError;
use aether_admin::provider::quota::parse_gemini_cli_retrieve_user_quota_response;
use aether_contracts::ProxySnapshot;
use aether_data_contracts::repository::provider_catalog::{
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
};
use aether_provider_pool::build_gemini_cli_pool_quota_request;
use serde_json::json;
use std::time::{SystemTime, UNIX_EPOCH};
async fn execute_gemini_cli_quota_plan(
state: &AdminAppState<'_>,
transport: &AdminGatewayProviderTransportSnapshot,
authorization: (String, String),
project_id: &str,
proxy_override: Option<&ProxySnapshot>,
) -> Result<ProviderQuotaExecutionOutcome, GatewayError> {
let proxy = match proxy_override {
Some(proxy) => Some(proxy.clone()),
None => {
state
.resolve_transport_proxy_snapshot_with_tunnel_affinity(transport)
.await
}
};
let timeouts = state
.resolve_transport_execution_timeouts(transport)
.or(Some(default_provider_quota_execution_timeouts(
proxy.as_ref(),
)));
let spec = build_gemini_cli_pool_quota_request(
&transport.key.id,
&transport.endpoint.base_url,
authorization,
project_id,
);
let plan = build_provider_quota_execution_plan(
transport,
spec,
proxy,
state.resolve_transport_profile(transport),
timeouts,
);
execute_provider_quota_plan(state, transport, plan, "gemini_cli").await
}
pub(crate) async fn refresh_gemini_cli_provider_quota_locally(
state: &AdminAppState<'_>,
provider: &StoredProviderCatalogProvider,
endpoint: &StoredProviderCatalogEndpoint,
keys: Vec<StoredProviderCatalogKey>,
proxy_override: Option<ProxySnapshot>,
) -> Result<Option<serde_json::Value>, GatewayError> {
let mut results = Vec::new();
let mut success_count = 0usize;
let mut failed_count = 0usize;
let mut auto_removed_count = 0usize;
for key in keys {
let mut transport = match state
.read_provider_transport_snapshot(&provider.id, &endpoint.id, &key.id)
.await?
{
Some(transport) => transport,
None => {
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": "Provider transport snapshot unavailable",
}));
continue;
}
};
let authorization = match state.resolve_local_oauth_header_auth(&transport).await? {
Some(auth) => auth,
_ => {
if quota_key_auto_removed(state, &key.id).await? {
auto_removed_count += 1;
results.push(oauth_refresh_auto_removed_result(&key));
continue;
}
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": "缺少 OAuth 认证信息,请先授权/刷新 Token",
}));
continue;
}
};
let project_id = match crate::provider_transport::resolve_gemini_cli_project_id(&transport)
{
Some(project_id) => Some(project_id),
None => state
.app()
.hydrate_gemini_cli_project_metadata_for_transport(&transport)
.await
.and_then(|hydrated| {
let project_id =
crate::provider_transport::resolve_gemini_cli_project_id(&hydrated);
transport = hydrated;
project_id
}),
};
let Some(project_id) = project_id else {
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": "缺少 Gemini CLI project_id,loadCodeAssist 未返回可用项目信息",
}));
continue;
};
let result = match execute_gemini_cli_quota_plan(
state,
&transport,
authorization,
&project_id,
proxy_override.as_ref(),
)
.await?
{
ProviderQuotaExecutionOutcome::Response(result) => result,
ProviderQuotaExecutionOutcome::Failure(detail) => {
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": format!("retrieveUserQuota 请求执行失败: {detail}"),
"status_code": 502,
}));
continue;
}
};
let now_unix_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()
.map(|duration| duration.as_secs())
.unwrap_or(0);
let mut metadata_update = None::<serde_json::Value>;
let (mut oauth_invalid_at_unix_secs, mut oauth_invalid_reason) =
quota_refresh_success_invalid_state(&key);
let mut status = "error".to_string();
let mut message = None::<String>;
if result.status_code == 200 {
if let Some(body_json) = result
.body
.as_ref()
.and_then(|body| body.json_body.as_ref())
{
metadata_update =
parse_gemini_cli_retrieve_user_quota_response(body_json, now_unix_secs)
.map(|metadata| json!({ "gemini_cli": metadata }));
if metadata_update.is_some() {
status = "success".to_string();
} else {
status = "no_metadata".to_string();
message = Some("响应中未包含配额 buckets".to_string());
}
} else {
status = "no_metadata".to_string();
message = Some("响应中未包含配额信息".to_string());
}
} else {
let err_msg = extract_execution_error_message(&result);
message = Some(match err_msg.as_deref() {
Some(detail) if !detail.is_empty() => {
format!(
"retrieveUserQuota 返回状态码 {}: {}",
result.status_code, detail
)
}
_ => format!("retrieveUserQuota 返回状态码 {}", result.status_code),
});
if result.status_code == 403 {
let reason = err_msg
.clone()
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| "账户访问被禁止".to_string());
oauth_invalid_at_unix_secs = Some(now_unix_secs);
oauth_invalid_reason = Some(format!("账户访问被禁止: {reason}"));
metadata_update = Some(json!({
"gemini_cli": {
"is_forbidden": true,
"forbidden_reason": reason,
"forbidden_at": now_unix_secs,
"updated_at": now_unix_secs,
}
}));
status = "forbidden".to_string();
}
}
if !persist_provider_quota_refresh_state(
state,
&key.id,
metadata_update.as_ref(),
oauth_invalid_at_unix_secs,
oauth_invalid_reason,
None,
)
.await?
{
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": "Key 状态写入失败",
}));
continue;
}
if status == "success" {
success_count += 1;
} else {
failed_count += 1;
}
let mut payload = serde_json::Map::new();
payload.insert("key_id".to_string(), json!(key.id));
payload.insert("key_name".to_string(), json!(key.name));
payload.insert("status".to_string(), json!(status));
if let Some(message) = message {
payload.insert("message".to_string(), json!(message));
}
if let Some(metadata) = metadata_update
.as_ref()
.and_then(|value| value.get("gemini_cli"))
.cloned()
{
payload.insert("metadata".to_string(), metadata);
}
if let Some(quota_snapshot) = build_quota_snapshot_payload(
"gemini_cli",
key.status_snapshot.as_ref(),
metadata_update.as_ref(),
) {
payload.insert("quota_snapshot".to_string(), quota_snapshot);
}
results.push(serde_json::Value::Object(payload));
}
Ok(Some(json!({
"success": success_count,
"failed": failed_count,
"total": results.len(),
"results": results,
"message": format!("已处理 {} 个 Key", results.len()),
"auto_removed": auto_removed_count,
})))
}
@@ -2,6 +2,7 @@ pub(crate) mod antigravity;
pub(crate) mod chatgpt_web;
pub(crate) mod codex;
pub(crate) mod dispatch;
pub(crate) mod gemini_cli;
pub(crate) mod grok;
pub(crate) mod kiro;
pub(crate) mod shared;
@@ -782,6 +782,18 @@ fn admin_pool_build_grok_account_quota_from_snapshot(
fn admin_pool_build_gemini_cli_account_quota_from_snapshot(
quota_snapshot: &serde_json::Map<String, serde_json::Value>,
) -> Option<String> {
if let Some(credits) = quota_snapshot
.get("credits")
.and_then(serde_json::Value::as_object)
{
if let Some(remaining) = admin_pool_json_to_f64(credits.get("remaining")) {
return Some(format!(
"AI Credits 剩余 {}",
admin_pool_format_quota_value(remaining)
));
}
}
let now = chrono::Utc::now().timestamp();
let mut active = admin_pool_quota_windows(quota_snapshot)
.into_iter()
@@ -15,7 +15,7 @@ use crate::ai_serving::{
ANTIGRAVITY_V1INTERNAL_ENVELOPE_NAME, GEMINI_CHAT_SYNC_FINALIZE_REPORT_KIND,
OPENAI_IMAGE_SYNC_FINALIZE_REPORT_KIND,
};
use crate::clock::current_unix_ms;
use crate::clock::{current_unix_ms, current_unix_secs};
use crate::execution_runtime;
use crate::handlers::admin::provider::shared::model_test_capabilities::{
admin_provider_model_supports_image_generation, admin_provider_model_test_capabilities_payload,
@@ -64,7 +64,7 @@ use aether_data_contracts::repository::provider_catalog::{
};
use aether_model_fetch::{
aggregate_models_for_cache, fetch_models_from_transports, json_string_list,
preset_models_for_provider, selected_models_fetch_endpoints,
merge_upstream_metadata, preset_models_for_provider, selected_models_fetch_endpoints,
};
use axum::{
body::{to_bytes, Body},
@@ -516,6 +516,18 @@ async fn provider_query_fetch_models_for_key(
)
.await;
}
if let Some(upstream_metadata) = outcome.upstream_metadata.as_ref() {
let merged_metadata =
merge_upstream_metadata(key.upstream_metadata.as_ref(), upstream_metadata);
state
.app()
.update_provider_catalog_key_upstream_metadata(
&key.id,
Some(&merged_metadata),
Some(current_unix_secs()),
)
.await?;
}
if unique_models.is_empty() && !all_errors.is_empty() {
if let Some(fallback) = provider_query_codex_preset_fallback(provider) {
@@ -1576,12 +1576,18 @@ fn provider_query_aggregate_standard_stream_sync_response(
fn provider_query_standard_execution_response_body(
provider_api_format: &str,
result: &aether_contracts::ExecutionResult,
report_context: Option<&Value>,
) -> Option<Value> {
let body = provider_query_execution_json_body(result).or_else(|| {
provider_query_decode_execution_body(result).and_then(|body| {
provider_query_aggregate_standard_stream_sync_response(provider_api_format, &body)
})
})?;
let body = report_context
.and_then(|context| {
crate::ai_serving::api::normalize_provider_private_response_value(body.clone(), context)
})
.unwrap_or(body);
if result.status_code < 400
&& provider_query_normalize_api_format_alias(provider_api_format)
== "gemini:generate_content"
@@ -2750,7 +2756,7 @@ async fn provider_query_execute_standard_test_candidate(
route_path: &str,
trace_id: &str,
) -> Result<ProviderQueryExecutionOutcome, GatewayError> {
let Some(transport) = state
let Some(mut transport) = state
.read_provider_transport_snapshot(&provider.id, &candidate.endpoint.id, &candidate.key.id)
.await?
else {
@@ -2963,6 +2969,56 @@ async fn provider_query_execute_standard_test_candidate(
upstream_is_stream,
require_body_stream_field,
);
if crate::provider_transport::is_gemini_cli_provider_transport(&transport)
&& normalized_provider_api_format == "gemini:generate_content"
{
let project_id = match crate::provider_transport::resolve_gemini_cli_project_id(&transport)
{
Some(project_id) => Some(project_id),
None => state
.app()
.hydrate_gemini_cli_project_metadata_for_transport(&transport)
.await
.and_then(|hydrated| {
let project_id =
crate::provider_transport::resolve_gemini_cli_project_id(&hydrated);
transport = hydrated;
project_id
}),
};
let Some(project_id) = project_id else {
return Ok(provider_query_skipped_execution_outcome(
provider_request_body,
"Gemini CLI project_id is unavailable for v1internal request",
));
};
provider_request_body = match crate::provider_transport::build_gemini_cli_v1internal_request(
project_id.as_str(),
trace_id,
request_model,
&provider_request_body,
) {
crate::provider_transport::GeminiCliRequestEnvelopeSupport::Supported(envelope) => {
envelope
}
crate::provider_transport::GeminiCliRequestEnvelopeSupport::Unsupported(_) => {
return Ok(provider_query_skipped_execution_outcome(
provider_request_body,
"Gemini CLI v1internal envelope could not be built",
));
}
};
}
let private_report_context =
(crate::provider_transport::is_gemini_cli_provider_transport(&transport)
&& normalized_provider_api_format == "gemini:generate_content")
.then(|| {
json!({
"has_envelope": true,
"envelope_name": crate::provider_transport::GEMINI_CLI_V1INTERNAL_ENVELOPE_NAME,
"provider_api_format": provider_api_format,
})
});
let uses_vertex_query_auth =
crate::provider_transport::uses_vertex_api_key_query_auth(&transport, provider_api_format);
@@ -3083,6 +3139,13 @@ async fn provider_query_execute_standard_test_candidate(
request_headers
.entry("content-type".to_string())
.or_insert_with(|| "application/json".to_string());
if crate::provider_transport::is_gemini_cli_provider_transport(&transport)
&& normalized_provider_api_format == "gemini:generate_content"
{
request_headers
.entry("user-agent".to_string())
.or_insert_with(|| crate::provider_transport::GEMINI_CLI_USER_AGENT.to_string());
}
let protected_headers = if uses_vertex_query_auth {
vec!["content-type"]
} else {
@@ -3160,7 +3223,11 @@ async fn provider_query_execute_standard_test_candidate(
.execute_execution_runtime_sync_plan(Some(trace_id), &plan)
.await?;
let response_body = if result.status_code < 400 {
provider_query_standard_execution_response_body(provider_api_format, &result)
provider_query_standard_execution_response_body(
provider_api_format,
&result,
private_report_context.as_ref(),
)
} else {
result.body.as_ref().and_then(|body| body.json_body.clone())
};
@@ -51,6 +51,7 @@ fn sample_openai_image_transport(provider_type: &str) -> AdminGatewayProviderTra
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: String::new(),
decrypted_auth_config: Some(
json!({
@@ -185,7 +186,7 @@ fn provider_query_execution_json_body_decodes_stream_encoded_json_response() {
Some(body.clone())
);
assert_eq!(
provider_query_standard_execution_response_body("openai:image", &result),
provider_query_standard_execution_response_body("openai:image", &result, None),
Some(body)
);
}
@@ -390,7 +391,7 @@ fn provider_query_standard_test_aggregates_responses_stream_body() {
error: None,
};
let body = provider_query_standard_execution_response_body("openai:responses", &result)
let body = provider_query_standard_execution_response_body("openai:responses", &result, None)
.expect("stream body should aggregate");
assert_eq!(body["model"], json!("gpt-5.4-mini"));
@@ -422,7 +423,7 @@ fn provider_query_standard_test_aggregates_responses_image_generation_call() {
error: None,
};
let body = provider_query_standard_execution_response_body("openai:responses", &result)
let body = provider_query_standard_execution_response_body("openai:responses", &result, None)
.expect("responses image stream body should aggregate");
assert_eq!(body["output"][0]["type"], json!("image_generation_call"));
@@ -584,10 +585,12 @@ fn provider_query_standard_test_rejects_gemini_success_without_visible_output()
error: None,
};
assert!(
provider_query_standard_execution_response_body("gemini:generate_content", &result)
.is_none()
);
assert!(provider_query_standard_execution_response_body(
"gemini:generate_content",
&result,
None
)
.is_none());
}
#[test]
@@ -295,6 +295,25 @@ impl<'a> AdminAppState<'a> {
use crate::handlers::public::{admin_requested_force_stream, normalize_admin_base_url};
use aether_admin::provider::endpoints as admin_provider_endpoints_pure;
let (fields, payload) = patch.into_parts();
let provider_type = provider.provider_type.trim().to_ascii_lowercase();
if provider_type == "gemini_cli"
&& [
"base_url",
"custom_path",
"header_rules",
"body_rules",
"max_retries",
"is_active",
"config",
"proxy",
"format_acceptance_config",
]
.iter()
.any(|field| fields.contains(field))
{
return Err("Gemini CLI Endpoint 由系统固定管理,不允许修改".to_string());
}
if self.provider_type_is_fixed(&provider.provider_type)
&& (fields.contains("base_url") || fields.contains("custom_path"))
@@ -326,7 +345,6 @@ impl<'a> AdminAppState<'a> {
&update_fields,
)?;
let provider_type = provider.provider_type.trim().to_ascii_lowercase();
if provider_type == "codex"
&& crate::ai_serving::is_openai_responses_format(&existing_endpoint.api_format)
{
@@ -580,6 +580,52 @@ fn provider_quota_metadata_string(
})
}
fn provider_quota_metadata_value_by_path<'a>(
metadata: &'a Map<String, Value>,
path: &[&str],
) -> Option<&'a Value> {
let (first, rest) = path.split_first()?;
let mut current = metadata.get(*first)?;
for segment in rest {
current = current.as_object()?.get(*segment)?;
}
Some(current)
}
fn provider_quota_metadata_number_by_paths(
metadata: &Map<String, Value>,
paths: &[&[&str]],
) -> Option<f64> {
paths.iter().find_map(|path| {
provider_quota_metadata_value_by_path(metadata, path)
.and_then(admin_provider_quota_pure::coerce_json_f64)
.filter(|value| value.is_finite())
})
}
fn provider_quota_metadata_bool_by_paths(
metadata: &Map<String, Value>,
paths: &[&[&str]],
) -> Option<bool> {
paths.iter().find_map(|path| {
provider_quota_metadata_value_by_path(metadata, path)
.and_then(admin_provider_quota_pure::coerce_json_bool)
})
}
fn provider_quota_metadata_string_by_paths(
metadata: &Map<String, Value>,
paths: &[&[&str]],
) -> Option<String> {
paths.iter().find_map(|path| {
provider_quota_metadata_value_by_path(metadata, path)
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
})
}
fn quota_windows_usage_ratio(windows: &[Value]) -> Option<f64> {
windows
.iter()
@@ -1531,12 +1577,144 @@ fn build_grok_quota_status_snapshot(
}))
}
fn gemini_cli_plan_type(metadata: &Map<String, Value>) -> Option<String> {
provider_quota_metadata_string(metadata, &["plan_type", "tier", "plan"]).or_else(|| {
provider_quota_metadata_string_by_paths(
metadata,
&[
&["paidTier", "id"],
&["paidTier", "tierType"],
&["paidTier", "name"],
&["currentTier", "id"],
&["currentTier", "tierType"],
&["currentTier", "name"],
],
)
})
}
fn gemini_cli_credits_status_snapshot(metadata: &Map<String, Value>) -> Option<Value> {
let remaining = provider_quota_metadata_number_by_paths(
metadata,
&[
&["credits", "remaining"],
&["credits", "remainingCredits"],
&["credits", "available"],
&["credits", "availableCredits"],
&["credits", "balance"],
&["remainingCredits"],
&["availableCredits"],
&["paidTier", "remainingCredits"],
&["paidTier", "availableCredits"],
&["currentTier", "remainingCredits"],
&["currentTier", "availableCredits"],
],
);
let balance = provider_quota_metadata_number_by_paths(
metadata,
&[
&["credits", "balance"],
&["credits", "remaining"],
&["credits", "available"],
&["paidTier", "availableCredits"],
&["currentTier", "availableCredits"],
&["availableCredits"],
&["remainingCredits"],
],
)
.or(remaining);
let consumed = provider_quota_metadata_number_by_paths(
metadata,
&[
&["credits", "consumed"],
&["credits", "consumedCredits"],
&["consumedCredits"],
&["paidTier", "consumedCredits"],
&["currentTier", "consumedCredits"],
],
);
let total = provider_quota_metadata_number_by_paths(
metadata,
&[
&["credits", "total"],
&["credits", "totalCredits"],
&["totalCredits"],
&["paidTier", "totalCredits"],
&["currentTier", "totalCredits"],
],
);
let unlimited = provider_quota_metadata_bool_by_paths(
metadata,
&[
&["credits", "unlimited"],
&["credits_unlimited"],
&["paidTier", "unlimited"],
&["currentTier", "unlimited"],
],
);
let explicit_has_credits = provider_quota_metadata_bool_by_paths(
metadata,
&[
&["credits", "has_credits"],
&["has_credits"],
&["paidTier", "hasCredits"],
&["currentTier", "hasCredits"],
],
);
let trace_id = provider_quota_metadata_string_by_paths(
metadata,
&[
&["credits", "trace_id"],
&["credits", "traceId"],
&["trace_id"],
&["traceId"],
],
);
let updated_at = provider_quota_metadata_value_by_path(metadata, &["credits", "updated_at"])
.or_else(|| provider_quota_metadata_value_by_path(metadata, &["credits", "updatedAt"]))
.or_else(|| metadata.get("updated_at"))
.and_then(|value| provider_quota_timestamp_unix_secs(Some(value)));
if remaining.is_none()
&& balance.is_none()
&& consumed.is_none()
&& total.is_none()
&& unlimited.is_none()
&& explicit_has_credits.is_none()
&& trace_id.is_none()
{
return None;
}
let has_credits = explicit_has_credits
.or_else(|| unlimited.filter(|value| *value))
.or_else(|| remaining.or(balance).map(|value| value > 0.0));
let mut credits = Map::new();
credits.insert("has_credits".to_string(), json!(has_credits));
credits.insert("balance".to_string(), json!(balance));
credits.insert("remaining".to_string(), json!(remaining.or(balance)));
credits.insert("consumed".to_string(), json!(consumed));
credits.insert("total".to_string(), json!(total));
credits.insert("unlimited".to_string(), json!(unlimited));
credits.insert("trace_id".to_string(), json!(trace_id));
credits.insert("updated_at".to_string(), json!(updated_at));
Some(Value::Object(credits))
}
fn build_gemini_cli_quota_status_snapshot(
upstream_metadata: Option<&Value>,
source: &str,
) -> Option<Value> {
let metadata = provider_quota_metadata_bucket(upstream_metadata, "gemini_cli")?;
let observed_at_unix_secs = provider_quota_timestamp_unix_secs(metadata.get("updated_at"));
let credits = gemini_cli_credits_status_snapshot(metadata);
let plan_type = gemini_cli_plan_type(metadata);
let observed_at_unix_secs = provider_quota_timestamp_unix_secs(metadata.get("updated_at"))
.or_else(|| {
credits
.as_ref()
.and_then(|value| value.get("updated_at"))
.and_then(|value| provider_quota_timestamp_unix_secs(Some(value)))
});
let windows = provider_quota_model_bucket(metadata)
.map(|models| {
models
@@ -1552,7 +1730,11 @@ fn build_gemini_cli_quota_status_snapshot(
})
.unwrap_or_default();
if windows.is_empty() && observed_at_unix_secs.is_none() {
if windows.is_empty()
&& observed_at_unix_secs.is_none()
&& credits.is_none()
&& plan_type.is_none()
{
return None;
}
@@ -1616,7 +1798,8 @@ fn build_gemini_cli_quota_status_snapshot(
"updated_at": observed_at_unix_secs,
"reset_at": reset_at,
"reset_seconds": reset_seconds,
"plan_type": serde_json::Value::Null,
"plan_type": plan_type,
"credits": credits,
"windows": windows,
}))
}
@@ -2683,6 +2866,55 @@ mod tests {
assert_eq!(auto.get("used_value"), Some(&json!(90.0)));
}
#[test]
fn provider_key_status_snapshot_payload_backfills_gemini_cli_account_credits() {
let mut key = sample_catalog_key();
key.upstream_metadata = Some(json!({
"gemini_cli": {
"updated_at": 1_778_067_246u64,
"plan_type": "g1-pro-tier",
"paidTier": {
"availableCredits": 123.5,
"consumedCredits": 7.0,
"totalCredits": 200.0
},
"quota_by_model": {
"gemini-2.5-pro": {
"display_name": "Gemini 2.5 Pro",
"remaining_fraction": 0.75,
"is_exhausted": false
}
}
}
}));
let payload = provider_key_status_snapshot_payload(&key, "gemini_cli");
let quota = payload
.get("quota")
.and_then(Value::as_object)
.expect("quota snapshot should be object");
let credits = quota
.get("credits")
.and_then(Value::as_object)
.expect("credits snapshot should exist");
let windows = quota
.get("windows")
.and_then(Value::as_array)
.expect("Gemini CLI model windows should exist");
assert_eq!(quota.get("provider_type"), Some(&json!("gemini_cli")));
assert_eq!(quota.get("code"), Some(&json!("ok")));
assert_eq!(quota.get("plan_type"), Some(&json!("g1-pro-tier")));
assert_eq!(quota.get("updated_at"), Some(&json!(1_778_067_246u64)));
assert_eq!(credits.get("remaining"), Some(&json!(123.5)));
assert_eq!(credits.get("balance"), Some(&json!(123.5)));
assert_eq!(credits.get("consumed"), Some(&json!(7.0)));
assert_eq!(credits.get("total"), Some(&json!(200.0)));
assert_eq!(credits.get("has_credits"), Some(&json!(true)));
assert_eq!(windows.len(), 1);
assert_eq!(windows[0].get("remaining_ratio"), Some(&json!(0.75)));
}
#[test]
fn provider_key_status_snapshot_payload_backfills_windsurf_daily_and_weekly_quota() {
let mut key = sample_catalog_key();
@@ -766,6 +766,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: decrypted_auth_config.map(ToOwned::to_owned),
},
@@ -251,6 +251,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -314,6 +314,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "secret".to_string(),
decrypted_auth_config: None,
},
@@ -9,6 +9,7 @@ use aether_usage_runtime::{
report_request_id, GatewayStreamReportRequest, GatewaySyncReportRequest,
GEMINI_FILE_MAPPING_TTL_SECONDS,
};
use base64::Engine as _;
use regex::Regex;
use serde_json::{json, Value};
use tracing::warn;
@@ -263,6 +264,121 @@ fn grok_upstream_response_body(report_context: Option<&Value>) -> Option<&Value>
.and_then(|response| response.get("body"))
}
fn gemini_cli_credits_from_report_context(
report_context: Option<&Value>,
now_unix_secs: u64,
) -> Option<Value> {
report_context
.and_then(|context| context.get("gemini_cli_v1internal_credits"))
.and_then(|value| {
admin_provider_quota_pure::parse_gemini_cli_v1internal_credits_response(
value,
now_unix_secs,
)
})
}
fn gemini_cli_credits_from_stream_payload(
payload: &GatewayStreamReportRequest,
now_unix_secs: u64,
) -> Option<Value> {
let body_base64 = payload.provider_body_base64.as_deref()?;
let body = base64::engine::general_purpose::STANDARD
.decode(body_base64)
.ok()?;
let text = std::str::from_utf8(&body).ok()?;
let mut latest = None::<Value>;
for raw_line in text.lines() {
let line = raw_line.trim_matches('\r').trim();
let data = line.strip_prefix("data:").map(str::trim).unwrap_or(line);
if data.is_empty() || data == "[DONE]" || data.starts_with(':') {
continue;
}
let Ok(value) = serde_json::from_str::<Value>(data) else {
continue;
};
if let Some(credits) =
admin_provider_quota_pure::parse_gemini_cli_v1internal_credits_response(
&value,
now_unix_secs,
)
{
latest = Some(credits);
}
}
latest
}
async fn sync_gemini_cli_credits_from_report(
state: &AppState,
report_context: Option<&Value>,
credits: Option<Value>,
) -> Result<bool, GatewayError> {
let Some(credits) = credits else {
return Ok(false);
};
let key_id = match report_context_key_id(report_context) {
Some(value) => value,
None => return Ok(false),
};
let Some(key) = state
.read_provider_catalog_keys_by_ids(std::slice::from_ref(&key_id))
.await?
.into_iter()
.next()
else {
return Ok(false);
};
let Some(provider) = state
.read_provider_catalog_providers_by_ids(std::slice::from_ref(&key.provider_id))
.await?
.into_iter()
.next()
else {
return Ok(false);
};
if !provider
.provider_type
.trim()
.eq_ignore_ascii_case("gemini_cli")
{
return Ok(false);
}
let now_unix_secs = current_unix_secs();
let mut gemini_cli_bucket = key
.upstream_metadata
.as_ref()
.and_then(Value::as_object)
.and_then(|metadata| metadata.get("gemini_cli"))
.and_then(Value::as_object)
.cloned()
.unwrap_or_else(serde_json::Map::new);
gemini_cli_bucket.insert("credits".to_string(), credits);
gemini_cli_bucket.insert("updated_at".to_string(), json!(now_unix_secs));
let updated_upstream_metadata = merge_metadata_object(
key.upstream_metadata.as_ref(),
"gemini_cli",
Value::Object(gemini_cli_bucket),
);
let updated_status_snapshot = sync_provider_key_quota_status_snapshot(
key.status_snapshot.as_ref(),
provider.provider_type.as_str(),
updated_upstream_metadata.as_ref(),
"report_effect",
);
let mut updated_key = key;
updated_key.upstream_metadata = updated_upstream_metadata;
updated_key.status_snapshot = updated_status_snapshot;
updated_key.updated_at_unix_secs = Some(now_unix_secs);
Ok(state
.update_provider_catalog_key(&updated_key)
.await?
.is_some())
}
fn grok_quota_reset_after_seconds(
body_json: Option<&Value>,
report_context: Option<&Value>,
@@ -464,6 +580,23 @@ async fn apply_local_sync_report_effect(state: &AppState, payload: &GatewaySyncR
"gateway failed to persist grok realtime quota from sync response"
);
}
let now_unix_secs = current_unix_secs();
if let Err(err) = sync_gemini_cli_credits_from_report(
state,
payload.report_context.as_ref(),
gemini_cli_credits_from_report_context(payload.report_context.as_ref(), now_unix_secs),
)
.await
{
warn!(
event_name = "gemini_cli_realtime_credits_sync_failed",
log_type = "ops",
report_kind = %payload.report_kind,
report_request_id = %short_request_id(report_request_id(payload.report_context.as_ref())),
error = ?err,
"gateway failed to persist gemini cli realtime credits from sync response"
);
}
}
async fn apply_local_stream_report_effect(state: &AppState, payload: &GatewayStreamReportRequest) {
@@ -500,6 +633,22 @@ async fn apply_local_stream_report_effect(state: &AppState, payload: &GatewayStr
"gateway failed to persist grok realtime quota from stream response"
);
}
let now_unix_secs = current_unix_secs();
let credits =
gemini_cli_credits_from_report_context(payload.report_context.as_ref(), now_unix_secs)
.or_else(|| gemini_cli_credits_from_stream_payload(payload, now_unix_secs));
if let Err(err) =
sync_gemini_cli_credits_from_report(state, payload.report_context.as_ref(), credits).await
{
warn!(
event_name = "gemini_cli_realtime_credits_sync_failed",
log_type = "ops",
report_kind = %payload.report_kind,
report_request_id = %short_request_id(report_request_id(payload.report_context.as_ref())),
error = ?err,
"gateway failed to persist gemini cli realtime credits from stream response"
);
}
}
async fn apply_local_gemini_file_mapping_report_effect(
+61 -3
View File
@@ -13,15 +13,16 @@ use aether_data_contracts::repository::provider_catalog::{
};
use aether_data_contracts::repository::quota::StoredProviderQuotaSnapshot;
use aether_model_fetch::{
aggregate_models_for_cache, model_fetch_interval_minutes, ModelFetchAssociationStore,
ModelFetchTransportRuntime,
aggregate_models_for_cache, fetch_models_from_transports, merge_upstream_metadata,
model_fetch_interval_minutes, ModelFetchAssociationStore, ModelFetchTransportRuntime,
};
use aether_scheduler_core::SchedulerAffinityTarget;
use async_trait::async_trait;
use serde_json::Value;
use tracing::debug;
use tracing::{debug, warn};
use super::{AppState, GatewayError};
use crate::clock::current_unix_secs;
use crate::model_fetch::ModelFetchRuntimeState;
use crate::provider_transport::{GatewayProviderTransportSnapshot, LocalResolvedOAuthRequestAuth};
use crate::request_candidate_runtime::{
@@ -31,6 +32,63 @@ use crate::request_candidate_runtime::{
use crate::scheduler::state::SchedulerRuntimeState;
use crate::{execution_runtime, provider_transport};
impl AppState {
pub(crate) async fn hydrate_gemini_cli_project_metadata_for_transport(
&self,
transport: &GatewayProviderTransportSnapshot,
) -> Option<GatewayProviderTransportSnapshot> {
if !provider_transport::is_gemini_cli_provider_transport(transport) {
return None;
}
if provider_transport::resolve_gemini_cli_project_id(transport).is_some() {
return Some(transport.clone());
}
let outcome =
match fetch_models_from_transports(self, std::slice::from_ref(transport)).await {
Ok(outcome) => outcome,
Err(err) => {
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
error = %err,
"gemini_cli project metadata hydration failed"
);
return None;
}
};
let upstream_metadata = outcome.upstream_metadata.as_ref()?;
let merged_metadata =
merge_upstream_metadata(transport.key.upstream_metadata.as_ref(), upstream_metadata);
let mut hydrated = transport.clone();
hydrated.key.upstream_metadata = Some(merged_metadata.clone());
if provider_transport::resolve_gemini_cli_project_id(&hydrated).is_none() {
return None;
}
if let Err(err) = self
.update_provider_catalog_key_upstream_metadata(
&transport.key.id,
Some(&merged_metadata),
Some(current_unix_secs()),
)
.await
{
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
error = ?err,
"gemini_cli project metadata hydration could not persist metadata"
);
}
Some(hydrated)
}
}
#[async_trait]
impl provider_transport::TransportTunnelAffinityLookup for AppState {
async fn lookup_tunnel_attachment_owner(
+1
View File
@@ -1811,6 +1811,7 @@ mod tests {
expires_at_unix_secs: None,
proxy: None,
fingerprint: None,
upstream_metadata: None,
decrypted_api_key: "__placeholder__".to_string(),
decrypted_auth_config: Some("{\"project_id\":\"demo\"}".to_string()),
},
@@ -453,6 +453,9 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
trace_id: String,
url: String,
has_model_field: bool,
project: String,
user_prompt_id: String,
envelope_model: String,
accept: String,
authorization: String,
exact_temperature: f64,
@@ -574,7 +577,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
)
.expect("endpoint should build")
.with_transport_fields(
"https://generativelanguage.googleapis.com".to_string(),
"https://cloudcode-pa.googleapis.com".to_string(),
Some(serde_json::json!([
{"action":"set","key":"x-endpoint-tag","value":"gemini-cli-oauth-local"}
])),
@@ -595,7 +598,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
fn sample_provider_catalog_key() -> StoredProviderCatalogKey {
let encrypted_auth_config = encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli","refresh_token":"rt-gemini-cli-stream-local-123"}"#,
r#"{"provider_type":"gemini_cli","refresh_token":"rt-gemini-cli-stream-local-123","project_id":"gemini-cli-project-1"}"#,
)
.expect("auth config should encrypt");
StoredProviderCatalogKey::new(
@@ -735,6 +738,27 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("model"))
.is_some(),
project: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("project"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
user_prompt_id: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("user_prompt_id"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
envelope_model: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("model"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
accept: payload
.get("headers")
.and_then(|value| value.get("accept"))
@@ -750,6 +774,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
exact_temperature: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("generationConfig"))
.and_then(|value| value.get("temperature"))
.and_then(|value| value.as_f64())
@@ -763,6 +788,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
metadata_mode: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("metadata"))
.and_then(|value| value.get("mode"))
.and_then(|value| value.as_str())
@@ -771,6 +797,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
metadata_source: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("metadata"))
.and_then(|value| value.get("source"))
.and_then(|value| value.as_str())
@@ -779,6 +806,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
tool_config_present: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("toolConfig"))
.is_some(),
proxy_node_id: payload
@@ -795,7 +823,7 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
});
let frames = concat!(
"{\"type\":\"headers\",\"payload\":{\"kind\":\"headers\",\"status_code\":200,\"headers\":{\"content-type\":\"text/event-stream\"}}}\n",
"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"data: {\\\"candidates\\\":[]}\\n\\n\"}}\n",
"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"data: {\\\"response\\\":{\\\"candidates\\\":[]},\\\"remainingCredits\\\":42,\\\"consumedCredits\\\":1,\\\"traceId\\\":\\\"trace-upstream-1\\\"}\\n\\n\"}}\n",
"{\"type\":\"telemetry\",\"payload\":{\"kind\":\"telemetry\",\"telemetry\":{\"elapsed_ms\":34,\"upstream_bytes\":26}}}\n",
"{\"type\":\"eof\",\"payload\":{\"kind\":\"eof\"}}\n"
);
@@ -909,9 +937,21 @@ async fn gateway_executes_gemini_cli_stream_via_local_decision_gate_after_oauth_
);
assert_eq!(
seen_execution_runtime_request.url,
"https://generativelanguage.googleapis.com/custom/v1beta/models/gemini-cli-upstream:streamGenerateContent?alt=sse"
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert!(seen_execution_runtime_request.has_model_field);
assert_eq!(
seen_execution_runtime_request.project,
"gemini-cli-project-1"
);
assert_eq!(
seen_execution_runtime_request.user_prompt_id,
"trace-gemini-cli-oauth-local-stream-123"
);
assert_eq!(
seen_execution_runtime_request.envelope_model,
"gemini-cli-upstream"
);
assert!(!seen_execution_runtime_request.has_model_field);
assert_eq!(seen_execution_runtime_request.accept, "text/event-stream");
assert_eq!(
seen_execution_runtime_request.authorization,
@@ -780,6 +780,9 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
trace_id: String,
url: String,
has_model_field: bool,
project: String,
user_prompt_id: String,
envelope_model: String,
authorization: String,
exact_temperature: f64,
endpoint_tag: String,
@@ -900,7 +903,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
)
.expect("endpoint should build")
.with_transport_fields(
"https://generativelanguage.googleapis.com".to_string(),
"https://cloudcode-pa.googleapis.com".to_string(),
Some(serde_json::json!([
{"action":"set","key":"x-endpoint-tag","value":"gemini-cli-oauth-local"}
])),
@@ -921,7 +924,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
fn sample_provider_catalog_key() -> StoredProviderCatalogKey {
let encrypted_auth_config = encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli","refresh_token":"rt-gemini-cli-local-123"}"#,
r#"{"provider_type":"gemini_cli","refresh_token":"rt-gemini-cli-local-123","project_id":"gemini-cli-project-1"}"#,
)
.expect("auth config should encrypt");
StoredProviderCatalogKey::new(
@@ -1062,6 +1065,27 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("model"))
.is_some(),
project: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("project"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
user_prompt_id: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("user_prompt_id"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
envelope_model: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("model"))
.and_then(|value| value.as_str())
.unwrap_or_default()
.to_string(),
authorization: payload
.get("headers")
.and_then(|value| value.get("authorization"))
@@ -1071,6 +1095,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
exact_temperature: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("generationConfig"))
.and_then(|value| value.get("temperature"))
.and_then(|value| value.as_f64())
@@ -1084,6 +1109,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
metadata_mode: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("metadata"))
.and_then(|value| value.get("mode"))
.and_then(|value| value.as_str())
@@ -1092,6 +1118,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
metadata_source: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("metadata"))
.and_then(|value| value.get("source"))
.and_then(|value| value.as_str())
@@ -1100,6 +1127,7 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
tool_config_present: payload
.get("body")
.and_then(|value| value.get("json_body"))
.and_then(|value| value.get("request"))
.and_then(|value| value.get("toolConfig"))
.is_some(),
proxy_node_id: payload
@@ -1123,18 +1151,23 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
},
"body": {
"json_body": {
"candidates": [{
"content": {
"role": "model",
"parts": [{"text": "Hello from Gemini CLI"}]
},
"finishReason": "STOP"
}],
"usageMetadata": {
"promptTokenCount": 1,
"candidatesTokenCount": 2,
"totalTokenCount": 3
"response": {
"candidates": [{
"content": {
"role": "model",
"parts": [{"text": "Hello from Gemini CLI"}]
},
"finishReason": "STOP"
}],
"usageMetadata": {
"promptTokenCount": 1,
"candidatesTokenCount": 2,
"totalTokenCount": 3
}
}
,"remainingCredits": 41,
"consumedCredits": 1,
"traceId": "trace-upstream-sync-1"
}
},
"telemetry": {
@@ -1203,7 +1236,9 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let response_status = response.status();
let response_body = response.text().await.expect("response body should read");
assert_eq!(response_status, StatusCode::OK, "body={response_body}");
let seen_refresh_request = seen_refresh
.lock()
@@ -1238,9 +1273,21 @@ async fn gateway_executes_gemini_cli_sync_via_local_decision_gate_after_oauth_re
);
assert_eq!(
seen_execution_runtime_request.url,
"https://generativelanguage.googleapis.com/custom/v1beta/models/gemini-cli-upstream:generateContent"
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert!(seen_execution_runtime_request.has_model_field);
assert_eq!(
seen_execution_runtime_request.project,
"gemini-cli-project-1"
);
assert_eq!(
seen_execution_runtime_request.user_prompt_id,
"trace-gemini-cli-oauth-local-sync-123"
);
assert_eq!(
seen_execution_runtime_request.envelope_model,
"gemini-cli-upstream"
);
assert!(!seen_execution_runtime_request.has_model_field);
assert_eq!(
seen_execution_runtime_request.authorization,
"Bearer refreshed-gemini-cli-access-token"
@@ -1782,6 +1782,7 @@ fn admin_provider_oauth_quota_mod_stays_thin() {
"refresh_codex_provider_quota_locally",
"refresh_kiro_provider_quota_locally",
"refresh_antigravity_provider_quota_locally",
"refresh_gemini_cli_provider_quota_locally",
"refresh_chatgpt_web_provider_quota_locally",
] {
assert!(
@@ -1358,6 +1358,209 @@ async fn gateway_refresh_kiro_quota_reconciles_missing_fixed_endpoint_before_ref
execution_runtime_handle.abort();
}
#[tokio::test]
async fn gateway_refreshes_admin_provider_quota_locally_for_gemini_cli_with_trusted_admin_principal(
) {
#[derive(Debug, Clone)]
struct SeenExecutionRuntimeRequest {
url: String,
authorization: String,
provider_api_format: String,
request_body: Option<serde_json::Value>,
}
let upstream_hits = Arc::new(Mutex::new(0usize));
let upstream_hits_clone = Arc::clone(&upstream_hits);
let upstream = Router::new().route(
"/api/admin/endpoints/providers/provider-gemini-cli/refresh-quota",
any(move |_request: Request| {
let upstream_hits_inner = Arc::clone(&upstream_hits_clone);
async move {
*upstream_hits_inner.lock().expect("mutex should lock") += 1;
(StatusCode::OK, Body::from("unexpected upstream hit"))
}
}),
);
let seen_execution_runtime = Arc::new(Mutex::new(None::<SeenExecutionRuntimeRequest>));
let seen_execution_runtime_clone = Arc::clone(&seen_execution_runtime);
let execution_runtime = Router::new().route(
"/v1/execute/sync",
any(move |request: Request| {
let seen_execution_runtime_inner = Arc::clone(&seen_execution_runtime_clone);
async move {
let plan: aether_contracts::ExecutionPlan = serde_json::from_slice(
&to_bytes(request.into_body(), usize::MAX)
.await
.expect("body should read"),
)
.expect("plan should parse");
*seen_execution_runtime_inner
.lock()
.expect("mutex should lock") = Some(SeenExecutionRuntimeRequest {
url: plan.url.clone(),
authorization: plan
.headers
.get("authorization")
.cloned()
.unwrap_or_default(),
provider_api_format: plan.provider_api_format.clone(),
request_body: plan.body.json_body.clone(),
});
let result = aether_contracts::ExecutionResult {
request_id: plan.request_id,
candidate_id: None,
status_code: 200,
headers: BTreeMap::new(),
body: Some(aether_contracts::ResponseBody {
json_body: Some(json!({
"buckets": [
{
"modelId": "gemini-2.5-pro",
"tokenType": "model",
"displayName": "Gemini 2.5 Pro",
"remainingFraction": 0.25,
"resetTime": "2030-01-01T00:00:00Z",
"isExhausted": false
},
{
"modelId": "gemini-2.5-flash",
"tokenType": "model",
"displayName": "Gemini 2.5 Flash",
"quotaInfo": {
"remainingFraction": 0.0,
"resetTime": "2030-01-01T01:00:00Z",
"isExhausted": true
}
}
]
})),
body_bytes_b64: None,
}),
telemetry: None,
error: None,
};
(StatusCode::OK, Json(result))
}
}),
);
let mut key = sample_key(
"key-gemini-cli-quota",
"provider-gemini-cli",
"gemini:generate_content",
"cached-gemini-cli-token",
);
key.auth_type = "oauth".to_string();
key.encrypted_auth_config = Some(
encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli","project_id":"gemini-cli-project-1"}"#,
)
.expect("auth config should encrypt"),
);
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![StoredProviderCatalogProvider::new(
"provider-gemini-cli".to_string(),
"gemini_cli".to_string(),
Some("https://example.com".to_string()),
"gemini_cli".to_string(),
)
.expect("provider should build")],
vec![sample_endpoint(
"endpoint-gemini-cli-quota",
"provider-gemini-cli",
"gemini:generate_content",
"https://cloudcode-pa.googleapis.com",
)],
vec![key],
));
let (_upstream_url, upstream_handle) = start_server(upstream).await;
let (execution_runtime_url, execution_runtime_handle) = start_server(execution_runtime).await;
let gateway = build_router_with_state(
build_state_with_execution_runtime_override(execution_runtime_url.clone())
.with_data_state_for_tests(
GatewayDataState::with_provider_catalog_repository_for_tests(
provider_catalog_repository.clone(),
)
.with_encryption_key_for_tests(DEVELOPMENT_ENCRYPTION_KEY),
),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = reqwest::Client::new()
.post(format!(
"{gateway_url}/api/admin/endpoints/providers/provider-gemini-cli/refresh-quota"
))
.header(GATEWAY_HEADER, "rust-phase3b")
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
.header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin")
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["success"], 1);
assert_eq!(payload["failed"], 0);
assert_eq!(payload["results"][0]["status"], "success");
assert_eq!(
payload["results"][0]["quota_snapshot"]["provider_type"],
"gemini_cli"
);
assert_eq!(
payload["results"][0]["quota_snapshot"]["windows"][0]["model"],
"gemini-2.5-pro"
);
let seen_request = seen_execution_runtime
.lock()
.expect("mutex should lock")
.clone()
.expect("execution runtime request should be captured");
assert_eq!(
seen_request.url,
"https://cloudcode-pa.googleapis.com/v1internal:retrieveUserQuota"
);
assert_eq!(seen_request.authorization, "Bearer cached-gemini-cli-token");
assert_eq!(
seen_request.provider_api_format,
"gemini_cli:retrieve_user_quota"
);
assert_eq!(
seen_request.request_body,
Some(json!({
"project": "gemini-cli-project-1"
}))
);
assert_eq!(*upstream_hits.lock().expect("mutex should lock"), 0);
let reloaded = provider_catalog_repository
.list_keys_by_ids(&["key-gemini-cli-quota".to_string()])
.await
.expect("keys should read");
assert_eq!(reloaded.len(), 1);
let upstream_metadata = reloaded[0]
.upstream_metadata
.as_ref()
.expect("upstream metadata should persist");
assert_eq!(
upstream_metadata["gemini_cli"]["quota_by_model"]["gemini-2.5-pro"]["remaining_fraction"],
json!(0.25)
);
assert_eq!(
upstream_metadata["gemini_cli"]["quota_by_model"]["gemini-2.5-flash"]["is_exhausted"],
json!(true)
);
gateway_handle.abort();
execution_runtime_handle.abort();
upstream_handle.abort();
}
#[tokio::test]
async fn gateway_refresh_quota_reconciles_unsupported_fixed_provider_endpoints_before_clear_message(
) {
@@ -1370,14 +1573,6 @@ async fn gateway_refresh_quota_reconciles_unsupported_fixed_provider_endpoints_b
"https://api.anthropic.com",
"Claude Code 暂不支持自动刷新额度",
),
(
"provider-gemini-cli-reconcile",
"gemini_cli",
1usize,
"gemini:generate_content",
"https://cloudcode-pa.googleapis.com",
"Gemini CLI 暂不支持自动刷新额度",
),
(
"provider-vertex-ai-reconcile",
"vertex_ai",
@@ -2085,6 +2085,99 @@ async fn gateway_renders_gemini_cli_account_quota_from_status_snapshot() {
assert_eq!(keys[0]["account_quota"], json!("Gemini 2.5 Pro 冷却中"));
}
#[tokio::test]
async fn gateway_prefers_gemini_cli_account_credits_over_model_quota_text() {
let mut provider = sample_provider("provider-gemini-cli", "gemini_cli", 10)
.with_transport_fields(
true,
false,
true,
None,
None,
None,
None,
None,
Some(json!({
"pool_advanced": {
"enabled": true,
"skip_exhausted_accounts": true
}
})),
);
provider.provider_type = "gemini_cli".to_string();
let mut key = sample_key(
"key-gemini-cli-credits",
"provider-gemini-cli",
"gemini:generate_content",
"oauth-placeholder",
);
key.name = "gemini cli credits".to_string();
key.auth_type = "oauth".to_string();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "gemini_cli",
"code": "ok",
"freshness": "fresh",
"source": "report_effect",
"observed_at": 1_775_553_285u64,
"exhausted": false,
"usage_ratio": 0.25,
"updated_at": 1_775_553_285u64,
"plan_type": "g1-pro-tier",
"credits": {
"remaining": 123.5,
"consumed": 7.0,
"has_credits": true
},
"windows": [
{
"code": "model:gemini-2.5-pro",
"label": "Gemini 2.5 Pro",
"scope": "model",
"unit": "percent",
"model": "gemini-2.5-pro",
"used_ratio": 0.25,
"remaining_ratio": 0.75,
"is_exhausted": false
}
]
}
}));
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
Vec::new(),
vec![key],
));
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
provider_catalog_repository,
));
let response = local_admin_pool_response(
&state,
http::Method::GET,
"/api/admin/pool/provider-gemini-cli/keys?page=1&page_size=50&status=all",
None,
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = serde_json::from_slice(
&to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read"),
)
.expect("json body should parse");
let keys = payload["keys"].as_array().expect("keys should be array");
assert_eq!(keys[0]["scheduling_status"], json!("available"));
assert_eq!(keys[0]["account_quota"], json!("AI Credits 剩余 123.5"));
}
#[tokio::test]
async fn gateway_formats_codex_quota_countdown_from_reset_after_seconds() {
let mut provider = sample_provider("provider-codex", "codex", 10).with_transport_fields(
@@ -8,7 +8,9 @@ use aether_data::repository::provider_catalog::InMemoryProviderCatalogReadReposi
use aether_data_contracts::repository::candidates::{
RequestCandidateReadRepository, RequestCandidateStatus,
};
use aether_data_contracts::repository::provider_catalog::StoredProviderCatalogEndpoint;
use aether_data_contracts::repository::provider_catalog::{
ProviderCatalogReadRepository, StoredProviderCatalogEndpoint,
};
use axum::body::Body;
use axum::routing::any;
use axum::{extract::Request, Json, Router};
@@ -5025,10 +5027,22 @@ async fn gateway_handles_gemini_cli_test_model_with_oauth_header_fallback() {
assert_eq!(plan.endpoint_id, "endpoint-gemini-cli");
assert_eq!(plan.key_id, "key-gemini-cli");
assert_eq!(plan.provider_api_format, "gemini:generate_content");
assert!(plan.stream);
assert_eq!(
plan.url,
"https://generativelanguage.googleapis.com/v1beta/models/gemini-2.5-pro:generateContent"
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["project"],
json!("project-1")
);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["model"],
json!("gemini-2.5-pro")
);
assert!(plan.body.json_body.as_ref().unwrap()["request"]
.get("contents")
.is_some());
assert_eq!(
plan.headers.get("authorization").map(String::as_str),
Some("Bearer cached-gemini-cli-token")
@@ -5072,7 +5086,7 @@ async fn gateway_handles_gemini_cli_test_model_with_oauth_header_fallback() {
key.encrypted_auth_config = Some(
aether_crypto::encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli"}"#,
r#"{"provider_type":"gemini_cli","project_id":"project-1"}"#,
)
.expect("auth config should encrypt"),
);
@@ -5082,7 +5096,7 @@ async fn gateway_handles_gemini_cli_test_model_with_oauth_header_fallback() {
"endpoint-gemini-cli",
"provider-gemini",
"gemini:generate_content",
"https://generativelanguage.googleapis.com",
"https://cloudcode-pa.googleapis.com",
)],
vec![key],
));
@@ -5123,6 +5137,176 @@ async fn gateway_handles_gemini_cli_test_model_with_oauth_header_fallback() {
execution_runtime_handle.abort();
}
#[tokio::test]
async fn gateway_hydrates_gemini_cli_project_id_from_load_code_assist_for_test_model() {
let seen_urls = Arc::new(Mutex::new(Vec::<String>::new()));
let seen_urls_clone = Arc::clone(&seen_urls);
let execution_runtime = Router::new().route(
"/v1/execute/sync",
any(move |Json(plan): Json<ExecutionPlan>| {
let seen_urls_inner = Arc::clone(&seen_urls_clone);
async move {
seen_urls_inner
.lock()
.expect("mutex should lock")
.push(plan.url.clone());
if plan.url == "https://cloudcode-pa.googleapis.com/v1internal:loadCodeAssist" {
assert_eq!(plan.model_name.as_deref(), Some("loadCodeAssist"));
assert_eq!(
plan.headers.get("authorization").map(String::as_str),
Some("Bearer cached-gemini-cli-token")
);
assert_eq!(
plan.body.json_body.as_ref().and_then(|body| body
.get("metadata")
.and_then(|metadata| metadata.get("pluginType"))),
Some(&json!("GEMINI"))
);
return Json(json!({
"request_id": plan.request_id,
"candidate_id": plan.candidate_id,
"status_code": 200,
"headers": {
"content-type": "application/json"
},
"body": {
"json_body": {
"cloudaicompanionProject": {
"id": "project-from-load-code-assist"
},
"currentTier": {
"id": "free"
}
}
}
}));
}
assert_eq!(
plan.url,
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert!(plan.stream);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["project"],
json!("project-from-load-code-assist")
);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["model"],
json!("gemini-2.5-pro")
);
Json(json!({
"request_id": plan.request_id,
"candidate_id": plan.candidate_id,
"status_code": 200,
"headers": {
"content-type": "application/json"
},
"body": {
"json_body": {
"id": "chatcmpl-gemini-cli-test-model",
"choices": [{
"message": {
"role": "assistant",
"content": "Hello from hydrated Gemini CLI"
}
}]
}
},
"telemetry": {
"elapsed_ms": 19
}
}))
}
}),
);
let (execution_runtime_url, execution_runtime_handle) = start_server(execution_runtime).await;
let mut provider = sample_provider("provider-gemini", "Gemini", 10);
provider.provider_type = "gemini_cli".to_string();
let mut key = sample_key(
"key-gemini-cli",
"provider-gemini",
"gemini:generate_content",
"cached-gemini-cli-token",
);
key.auth_type = "oauth".to_string();
key.encrypted_auth_config = Some(
aether_crypto::encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli","refresh_token":"rt-gemini-cli-123"}"#,
)
.expect("auth config should encrypt"),
);
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
vec![sample_endpoint(
"endpoint-gemini-cli",
"provider-gemini",
"gemini:generate_content",
"https://cloudcode-pa.googleapis.com",
)],
vec![key],
));
let gateway = build_router_with_state(
build_state_with_execution_runtime_override(execution_runtime_url)
.with_data_state_for_tests(
GatewayDataState::with_provider_catalog_repository_for_tests(Arc::clone(
&provider_catalog_repository,
))
.with_encryption_key_for_tests(DEVELOPMENT_ENCRYPTION_KEY),
),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = reqwest::Client::new()
.post(format!("{gateway_url}/api/admin/provider-query/test-model"))
.header(GATEWAY_HEADER, "rust-phase3b")
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
.header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin")
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
.json(&json!({
"provider_id": "provider-gemini",
"model": "gemini-2.5-pro",
"api_format": "gemini:generate_content"
}))
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["success"], json!(true));
assert_eq!(
payload["data"]["response"]["choices"][0]["message"]["content"],
json!("Hello from hydrated Gemini CLI")
);
assert_eq!(
*seen_urls.lock().expect("mutex should lock"),
vec![
"https://cloudcode-pa.googleapis.com/v1internal:loadCodeAssist".to_string(),
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
.to_string(),
]
);
let reloaded = provider_catalog_repository
.list_keys_by_ids(&["key-gemini-cli".to_string()])
.await
.expect("key should reload");
assert_eq!(
reloaded[0]
.upstream_metadata
.as_ref()
.and_then(|metadata| metadata.get("gemini_cli"))
.and_then(|metadata| metadata.get("project_id")),
Some(&json!("project-from-load-code-assist"))
);
gateway_handle.abort();
execution_runtime_handle.abort();
}
#[tokio::test]
async fn gateway_uses_compatible_gemini_cli_endpoint_when_api_format_is_omitted() {
let execution_runtime = Router::new().route(
@@ -5303,6 +5487,118 @@ async fn gateway_handles_gemini_cli_test_model_failover_locally() {
execution_runtime_handle.abort();
}
#[tokio::test]
async fn gateway_unwraps_gemini_cli_v1internal_response_for_failover_model_test() {
let execution_runtime = Router::new().route(
"/v1/execute/sync",
any(move |Json(plan): Json<ExecutionPlan>| async move {
assert_eq!(plan.provider_id, "provider-gemini-cli");
assert_eq!(plan.endpoint_id, "endpoint-gemini-cli");
assert_eq!(plan.key_id, "key-gemini-cli");
assert_eq!(plan.provider_api_format, "gemini:generate_content");
assert!(plan.stream);
assert_eq!(
plan.url,
"https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse"
);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["project"],
json!("project-1")
);
assert_eq!(
plan.body.json_body.as_ref().unwrap()["model"],
json!("gemini-3-flash-preview")
);
Json(json!({
"request_id": plan.request_id,
"candidate_id": plan.candidate_id,
"status_code": 200,
"headers": {
"content-type": "text/event-stream"
},
"body": {
"body_bytes_b64": base64::engine::general_purpose::STANDARD.encode(
concat!(
"data: {\"response\":{\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"Gemini CLI v1internal failover response\"}],\"role\":\"model\"},\"finishReason\":\"STOP\",\"index\":0}],\"modelVersion\":\"gemini-3-flash-preview\",\"usageMetadata\":{\"promptTokenCount\":2,\"candidatesTokenCount\":5,\"totalTokenCount\":7}},\"remainingCredits\":123,\"consumedCredits\":1,\"traceId\":\"trace-gemini-cli-1\"}\n\n"
)
.as_bytes()
)
},
"telemetry": {
"elapsed_ms": 23
}
}))
}),
);
let (execution_runtime_url, execution_runtime_handle) = start_server(execution_runtime).await;
let mut provider = sample_provider("provider-gemini-cli", "Gemini CLI", 10);
provider.provider_type = "gemini_cli".to_string();
let mut key = sample_key(
"key-gemini-cli",
"provider-gemini-cli",
"gemini:generate_content",
"cached-gemini-cli-token",
);
key.auth_type = "oauth".to_string();
key.encrypted_auth_config = Some(
aether_crypto::encrypt_python_fernet_plaintext(
DEVELOPMENT_ENCRYPTION_KEY,
r#"{"provider_type":"gemini_cli","project_id":"project-1"}"#,
)
.expect("auth config should encrypt"),
);
let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
vec![sample_endpoint(
"endpoint-gemini-cli",
"provider-gemini-cli",
"gemini:generate_content",
"https://cloudcode-pa.googleapis.com",
)],
vec![key],
));
let gateway = build_router_with_state(
build_state_with_execution_runtime_override(execution_runtime_url)
.with_data_state_for_tests(GatewayDataState::with_provider_transport_reader_for_tests(
provider_catalog_repository,
DEVELOPMENT_ENCRYPTION_KEY.to_string(),
)),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = reqwest::Client::new()
.post(format!(
"{gateway_url}/api/admin/provider-query/test-model-failover"
))
.header(GATEWAY_HEADER, "rust-phase3b")
.header(TRUSTED_ADMIN_USER_ID_HEADER, "admin-user-123")
.header(TRUSTED_ADMIN_USER_ROLE_HEADER, "admin")
.header(TRUSTED_ADMIN_SESSION_ID_HEADER, "session-123")
.json(&json!({
"provider_id": "provider-gemini-cli",
"failover_models": ["gemini-3-flash-preview"],
"api_format": "gemini:generate_content"
}))
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["success"], json!(true));
assert_eq!(payload["total_attempts"], json!(1));
assert_eq!(
payload["data"]["response"]["candidates"][0]["content"]["parts"][0]["text"],
json!("Gemini CLI v1internal failover response")
);
assert!(payload["data"]["response"].get("response").is_none());
gateway_handle.abort();
execution_runtime_handle.abort();
}
#[tokio::test]
async fn gateway_handles_admin_provider_query_test_model_failover_with_single_model_name_alias() {
let execution_runtime = Router::new().route(