修复 Antigravity OAuth 配额复检缺 project

This commit is contained in:
MMEXA
2026-06-19 22:21:22 +08:00
parent 16584067d7
commit 6c4e730e60
12 changed files with 977 additions and 130 deletions
@@ -14,6 +14,7 @@ use crate::ai_serving::transport::antigravity::{
build_antigravity_safe_v1internal_request, build_antigravity_static_identity_headers,
classify_local_antigravity_request_support, AntigravityEnvelopeRequestType,
AntigravityRequestEnvelopeSupport, AntigravityRequestSideSupport,
AntigravityRequestSideUnsupportedReason,
};
use crate::ai_serving::transport::{
build_gemini_cli_v1internal_request, build_grok_browser_headers, build_grok_upstream_url,
@@ -230,11 +231,32 @@ 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(
let mut antigravity_support = classify_local_antigravity_request_support(
&transport,
&base_provider_request_body,
AntigravityEnvelopeRequestType::Agent,
);
if matches!(
antigravity_support,
AntigravityRequestSideSupport::Unsupported(
AntigravityRequestSideUnsupportedReason::UnsupportedAuth(
crate::provider_transport::antigravity::AntigravityRequestAuthUnsupportedReason::MissingProjectId
)
)
) {
if let Some(hydrated) = state
.hydrate_antigravity_project_metadata_for_transport(&transport)
.await
{
transport = Arc::new(hydrated);
antigravity_support = classify_local_antigravity_request_support(
&transport,
&base_provider_request_body,
AntigravityEnvelopeRequestType::Agent,
);
}
}
match antigravity_support {
AntigravityRequestSideSupport::Supported(spec) => Some(spec.auth),
AntigravityRequestSideSupport::Unsupported(_) => {
mark_skipped_local_same_format_provider_candidate(
@@ -33,7 +33,7 @@ use crate::ai_serving::transport::antigravity::{
build_antigravity_safe_v1internal_request, build_antigravity_static_identity_headers,
classify_local_antigravity_request_support, is_antigravity_provider_transport,
AntigravityEnvelopeRequestType, AntigravityRequestEnvelopeSupport,
AntigravityRequestSideSupport,
AntigravityRequestSideSupport, AntigravityRequestSideUnsupportedReason,
};
use crate::ai_serving::transport::auth::{
resolve_local_gemini_auth, resolve_local_openai_bearer_auth, resolve_local_standard_auth,
@@ -119,11 +119,11 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
let provider_api_format = eligible.provider_api_format.as_str();
let normalized_provider_api_format =
crate::ai_serving::normalize_api_format_alias(provider_api_format);
let transport = &eligible.transport;
let transport_profile = crate::ai_serving::transport::resolve_transport_profile(transport);
let is_antigravity = is_antigravity_provider_transport(transport);
let is_gemini_cli = is_gemini_cli_provider_transport(transport);
let is_kiro_claude_cli = is_kiro_claude_messages_transport(transport, provider_api_format);
let mut transport = Arc::clone(&eligible.transport);
let transport_profile = crate::ai_serving::transport::resolve_transport_profile(&transport);
let is_antigravity = is_antigravity_provider_transport(&transport);
let is_gemini_cli = is_gemini_cli_provider_transport(&transport);
let is_kiro_claude_cli = is_kiro_claude_messages_transport(&transport, provider_api_format);
let is_grok = transport
.provider
.provider_type
@@ -145,33 +145,34 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
.await);
}
let is_windsurf_cascade =
provider_api_format == "openai:chat" && is_windsurf_provider_transport(transport);
provider_api_format == "openai:chat" && is_windsurf_provider_transport(&transport);
let same_format = api_format_alias_matches(provider_api_format, &client_api_format);
let conversion_kind = request_conversion_kind(spec_metadata.api_format, provider_api_format);
let transport_unsupported_reason =
if is_grok && is_grok_text_provider_api_format(provider_api_format) {
None
} else if same_format && is_kiro_claude_cli {
local_kiro_request_transport_unsupported_reason_with_network(transport)
} else if same_format {
local_standard_transport_unsupported_reason_with_network(transport, provider_api_format)
} else if is_windsurf_cascade {
local_windsurf_request_transport_unsupported_reason_with_network(transport)
} else {
match conversion_kind {
Some(_)
if (is_antigravity || is_gemini_cli)
&& normalized_provider_api_format == "gemini:generate_content" =>
{
None
}
Some(kind) => crate::ai_serving::request_conversion_transport_unsupported_reason(
transport, kind,
),
None => Some("transport_api_format_unsupported"),
let transport_unsupported_reason = if is_grok
&& is_grok_text_provider_api_format(provider_api_format)
{
None
} else if same_format && is_kiro_claude_cli {
local_kiro_request_transport_unsupported_reason_with_network(&transport)
} else if same_format {
local_standard_transport_unsupported_reason_with_network(&transport, provider_api_format)
} else if is_windsurf_cascade {
local_windsurf_request_transport_unsupported_reason_with_network(&transport)
} else {
match conversion_kind {
Some(_)
if (is_antigravity || is_gemini_cli)
&& normalized_provider_api_format == "gemini:generate_content" =>
{
None
}
};
Some(kind) => {
crate::ai_serving::request_conversion_transport_unsupported_reason(&transport, kind)
}
None => Some("transport_api_format_unsupported"),
}
};
if let Some(skip_reason) = transport_unsupported_reason {
mark_skipped_local_openai_responses_candidate(
state,
@@ -194,7 +195,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
let kiro_auth = if is_kiro_claude_cli {
match crate::ai_serving::planner::candidate_preparation::resolve_candidate_oauth_auth(
planner_state,
transport,
&transport,
oauth_context,
)
.await
@@ -219,20 +220,20 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
};
let direct_auth = if is_grok && is_grok_text_provider_api_format(provider_api_format) {
crate::ai_serving::transport::resolve_grok_session_auth(transport)
crate::ai_serving::transport::resolve_grok_session_auth(&transport)
} else if kiro_auth.is_some() {
None
} else if same_format {
match crate::ai_serving::normalize_api_format_alias(provider_api_format).as_str() {
"gemini:generate_content" => resolve_local_gemini_auth(transport),
"claude:messages" => resolve_local_standard_auth(transport),
"gemini:generate_content" => resolve_local_gemini_auth(&transport),
"claude:messages" => resolve_local_standard_auth(&transport),
"openai:responses" | "openai:responses:compact" => {
resolve_local_openai_bearer_auth(transport)
resolve_local_openai_bearer_auth(&transport)
}
_ => None,
}
} else {
conversion_kind.and_then(|kind| request_conversion_direct_auth(transport, kind))
conversion_kind.and_then(|kind| request_conversion_direct_auth(&transport, kind))
};
let prepared_candidate = if let Some(kiro_auth) = kiro_auth.as_ref() {
match prepare_header_authenticated_candidate_from_auth(
@@ -258,7 +259,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
} else {
match prepare_header_authenticated_candidate(
planner_state,
transport,
&transport,
candidate,
direct_auth,
oauth_context,
@@ -414,11 +415,32 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
Some(body_json),
);
let antigravity_auth = if is_antigravity {
match classify_local_antigravity_request_support(
transport,
let mut antigravity_support = classify_local_antigravity_request_support(
&transport,
&base_provider_request_body,
AntigravityEnvelopeRequestType::Agent,
);
if matches!(
antigravity_support,
AntigravityRequestSideSupport::Unsupported(
AntigravityRequestSideUnsupportedReason::UnsupportedAuth(
crate::provider_transport::antigravity::AntigravityRequestAuthUnsupportedReason::MissingProjectId
)
)
) {
if let Some(hydrated) = state
.hydrate_antigravity_project_metadata_for_transport(&transport)
.await
{
transport = Arc::new(hydrated);
antigravity_support = classify_local_antigravity_request_support(
&transport,
&base_provider_request_body,
AntigravityEnvelopeRequestType::Agent,
);
}
}
match antigravity_support {
AntigravityRequestSideSupport::Supported(spec) => Some(spec.auth),
AntigravityRequestSideSupport::Unsupported(_) => {
mark_skipped_local_openai_responses_candidate(
@@ -480,7 +502,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
candidate_index,
candidate_id,
spec_metadata.api_format,
transport,
&transport,
provider_api_format,
mapped_model,
auth_header,
@@ -504,7 +526,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
candidate_index,
candidate_id,
spec_metadata.api_format,
transport,
&transport,
provider_api_format,
mapped_model,
auth_header,
@@ -516,7 +538,7 @@ 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)
&& is_gemini_cli_provider_transport(&transport)
{
return Ok(build_gemini_cli_openai_responses_payload_parts(
state,
@@ -528,7 +550,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
candidate_index,
candidate_id,
spec_metadata.api_format,
transport,
&transport,
provider_api_format,
mapped_model,
auth_header,
@@ -541,11 +563,11 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
}
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))
Some(build_grok_upstream_url(&transport, GROK_CHAT_PATH))
} else if needs_bidirectional_conversion {
build_cross_format_openai_responses_upstream_url(
parts,
transport,
&transport,
&mapped_model,
spec_metadata.api_format,
provider_api_format,
@@ -554,7 +576,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
} else {
build_local_openai_responses_upstream_url(
parts,
transport,
&transport,
api_format_alias_matches(provider_api_format, "openai:responses:compact"),
)
}) else {
@@ -581,7 +603,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
.unwrap_or_default();
let resolved_headers = if is_grok && is_grok_text_provider_api_format(provider_api_format) {
let Some(headers) = build_grok_browser_headers(GrokHeaderInput {
transport,
transport: &transport,
transport_profile: transport_profile.as_ref(),
request_headers: Some(effective_headers),
content_type: "application/json",
@@ -615,7 +637,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
} else {
let Some(resolved_headers) =
build_standard_provider_request_headers(StandardProviderRequestHeadersInput {
transport,
transport: &transport,
provider_api_format,
same_format,
headers: effective_headers,
@@ -709,7 +731,7 @@ pub(crate) async fn resolve_local_openai_responses_candidate_payload_parts(
None
},
upstream_is_stream,
transport: Arc::clone(transport),
transport: Arc::clone(&transport),
transport_profile,
image_request_summary: None,
request_redacted: redaction.redacted,
@@ -32,6 +32,8 @@ pub(super) struct AdminProviderOAuthBatchImportEntry {
pub email: Option<String>,
pub account_name: Option<String>,
pub project_id: Option<String>,
pub client_version: Option<String>,
pub session_id: Option<String>,
pub sso_rw_token: Option<String>,
pub cf_cookies: Option<String>,
pub cf_clearance: Option<String>,
@@ -191,6 +193,8 @@ fn extract_admin_provider_oauth_batch_import_entry(
email: None,
account_name: None,
project_id: None,
client_version: None,
session_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"),
@@ -331,6 +335,19 @@ fn extract_admin_provider_oauth_batch_import_entry(
.or_else(|| object.get("cloudaicompanionProject"))
.or_else(|| object.get("cloudAiCompanionProject")),
);
let client_version = coerce_admin_provider_oauth_import_str(
object
.get("client_version")
.or_else(|| object.get("clientVersion"))
.or_else(|| object.get("antigravityClientVersion")),
);
let session_id = coerce_admin_provider_oauth_import_str(
object
.get("session_id")
.or_else(|| object.get("sessionId"))
.or_else(|| object.get("vscode_session_id"))
.or_else(|| object.get("vscodeSessionId")),
);
let sso_rw_token = coerce_admin_provider_oauth_import_str(
object
.get("sso_rw_token")
@@ -379,6 +396,8 @@ fn extract_admin_provider_oauth_batch_import_entry(
email,
account_name,
project_id,
client_version,
session_id,
sso_rw_token,
cf_cookies,
cf_clearance,
@@ -470,6 +489,8 @@ fn parse_error_entry(error: String) -> AdminProviderOAuthBatchImportEntry {
email: None,
account_name: None,
project_id: None,
client_version: None,
session_id: None,
sso_rw_token: None,
cf_cookies: None,
cf_clearance: None,
@@ -502,6 +523,29 @@ pub(super) fn apply_admin_provider_oauth_batch_import_hints(
}
return;
}
if provider_type == "antigravity" {
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(client_version) = entry.client_version.as_ref() {
auth_config
.entry("client_version".to_string())
.or_insert_with(|| json!(client_version));
}
if let Some(session_id) = entry.session_id.as_ref() {
auth_config
.entry("session_id".to_string())
.or_insert_with(|| json!(session_id));
}
if let Some(user_agent) = entry.user_agent.as_ref() {
auth_config
.entry("user_agent".to_string())
.or_insert_with(|| json!(user_agent));
}
return;
}
if !matches!(provider_type.as_str(), "codex" | "chatgpt_web" | "grok") {
return;
}
@@ -823,6 +867,23 @@ mod tests {
);
}
#[test]
fn applies_antigravity_project_and_user_agent_hints_to_auth_config() {
let entries = parse_admin_provider_oauth_batch_import_entries(
"antigravity",
r#"{"refreshToken":"rt-1","cloudaicompanionProject":{"id":"project-antigravity-2"},"userAgent":"antigravity"}"#,
);
let mut auth_config = serde_json::Map::new();
apply_admin_provider_oauth_batch_import_hints("antigravity", &entries[0], &mut auth_config);
assert_eq!(
auth_config.get("project_id"),
Some(&json!("project-antigravity-2"))
);
assert_eq!(auth_config.get("user_agent"), Some(&json!("antigravity")));
}
#[test]
fn parses_windsurf_json_credentials_for_native_import() {
let entries = parse_admin_provider_oauth_batch_import_entries(
@@ -106,6 +106,34 @@ fn import_payload_string_any(
.map(ToOwned::to_owned)
}
fn import_payload_project_id_any(
payload: &serde_json::Map<String, serde_json::Value>,
keys: &[&str],
) -> Option<String> {
keys.iter().find_map(|key| {
let value = payload.get(*key)?;
if let Some(string) = value
.as_str()
.map(str::trim)
.filter(|value| !value.is_empty())
{
return Some(string.to_string());
}
value
.as_object()
.and_then(|object| {
object
.get("id")
.or_else(|| object.get("project_id"))
.or_else(|| object.get("projectId"))
})
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
})
}
fn import_payload_u64_any(
payload: &serde_json::Map<String, serde_json::Value>,
keys: &[&str],
@@ -129,6 +157,49 @@ fn apply_single_import_hints(
auth_config: &mut serde_json::Map<String, serde_json::Value>,
) {
let provider_type = provider_type.trim().to_ascii_lowercase();
if provider_type == "antigravity" {
if let Some(project_id) = import_payload_project_id_any(
payload,
&[
"project_id",
"projectId",
"cloudaicompanionProject",
"cloudAiCompanionProject",
],
) {
auth_config
.entry("project_id".to_string())
.or_insert_with(|| json!(project_id));
}
for (target, keys) in [
(
"client_version",
&[
"client_version",
"clientVersion",
"antigravityClientVersion",
][..],
),
(
"session_id",
&[
"session_id",
"sessionId",
"vscode_session_id",
"vscodeSessionId",
][..],
),
("user_agent", &["user_agent", "userAgent"][..]),
] {
let Some(value) = import_payload_string_any(payload, keys) else {
continue;
};
auth_config
.entry(target.to_string())
.or_insert_with(|| json!(value));
}
return;
}
if !matches!(provider_type.as_str(), "codex" | "chatgpt_web" | "grok") {
return;
}
@@ -640,7 +711,8 @@ pub(super) async fn handle_admin_provider_oauth_import_refresh_token(
#[cfg(test)]
mod tests {
use super::{
import_payload_string_any, import_payload_u64_any, sanitize_windsurf_import_error,
apply_single_import_hints, import_payload_string_any, import_payload_u64_any,
sanitize_windsurf_import_error,
};
use aether_oauth::core::OAuthError;
use serde_json::json;
@@ -675,6 +747,35 @@ mod tests {
);
}
#[test]
fn single_import_applies_antigravity_identity_hints() {
let payload = json!({
"cloudaicompanionProject": {
"id": "project-antigravity-1"
},
"clientVersion": "1.99.0",
"sessionId": "session-antigravity-1",
"userAgent": "antigravity"
})
.as_object()
.cloned()
.expect("payload should be an object");
let mut auth_config = serde_json::Map::new();
apply_single_import_hints("antigravity", &payload, &mut auth_config);
assert_eq!(
auth_config.get("project_id"),
Some(&json!("project-antigravity-1"))
);
assert_eq!(auth_config.get("client_version"), Some(&json!("1.99.0")));
assert_eq!(
auth_config.get("session_id"),
Some(&json!("session-antigravity-1"))
);
assert_eq!(auth_config.get("user_agent"), Some(&json!("antigravity")));
}
#[test]
fn windsurf_import_error_redacts_http_body() {
let error = OAuthError::HttpStatus {
@@ -68,7 +68,7 @@ pub(crate) async fn refresh_antigravity_provider_quota_locally(
let mut auto_removed_count = 0usize;
for key in keys {
let transport = match state
let mut transport = match state
.read_provider_transport_snapshot(&provider.id, &endpoint.id, &key.id)
.await?
{
@@ -104,15 +104,25 @@ pub(crate) async fn refresh_antigravity_provider_quota_locally(
}
};
let Some((project_id, identity_headers)) =
state.resolve_local_antigravity_identity_headers(&transport)
else {
let identity = match state.resolve_local_antigravity_identity_headers(&transport) {
Some(identity) => Some(identity),
None => state
.app()
.hydrate_antigravity_project_metadata_for_transport(&transport)
.await
.and_then(|hydrated| {
let identity = state.resolve_local_antigravity_identity_headers(&hydrated);
transport = hydrated;
identity
}),
};
let Some((project_id, identity_headers)) = identity else {
failed_count += 1;
results.push(json!({
"key_id": key.id,
"key_name": key.name,
"status": "error",
"message": "缺少 OAuth 认证信息,请先授权/刷新 Token",
"message": "缺少 Antigravity project_id,loadCodeAssist 未返回可用项目信息",
}));
continue;
};
+130 -2
View File
@@ -13,8 +13,9 @@ use aether_data_contracts::repository::provider_catalog::{
};
use aether_data_contracts::repository::quota::StoredProviderQuotaSnapshot;
use aether_model_fetch::{
aggregate_models_for_cache, fetch_models_from_transports, merge_upstream_metadata,
model_fetch_interval_minutes, ModelFetchAssociationStore, ModelFetchTransportRuntime,
aggregate_models_for_cache, build_antigravity_load_code_assist_plan,
fetch_models_from_transports, merge_upstream_metadata, model_fetch_interval_minutes,
ModelFetchAssociationStore, ModelFetchTransportRuntime,
};
use aether_scheduler_core::SchedulerAffinityTarget;
use async_trait::async_trait;
@@ -33,6 +34,109 @@ use crate::scheduler::state::SchedulerRuntimeState;
use crate::{execution_runtime, provider_transport};
impl AppState {
pub(crate) async fn hydrate_antigravity_project_metadata_for_transport(
&self,
transport: &GatewayProviderTransportSnapshot,
) -> Option<GatewayProviderTransportSnapshot> {
if !provider_transport::antigravity::is_antigravity_provider_transport(transport) {
return None;
}
if matches!(
provider_transport::antigravity::resolve_local_antigravity_request_auth(transport),
provider_transport::antigravity::AntigravityRequestAuthSupport::Supported(_)
) {
return Some(transport.clone());
}
let plan = match build_antigravity_load_code_assist_plan(self, transport).await {
Ok(plan) => plan,
Err(err) => {
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
error = %err,
"antigravity project metadata hydration failed"
);
return None;
}
};
let result =
match execution_runtime::execute_execution_runtime_sync_plan(self, None, &plan).await {
Ok(result) => result,
Err(err) => {
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
error = ?err,
"antigravity project metadata hydration request failed"
);
return None;
}
};
if !(200..300).contains(&result.status_code) {
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
status_code = result.status_code,
"antigravity project metadata hydration returned non-success status"
);
return None;
}
let Some(project_id) = result
.body
.as_ref()
.and_then(|body| body.json_body.as_ref())
.and_then(extract_antigravity_load_code_assist_project_id)
else {
warn!(
provider_id = %transport.provider.id,
endpoint_id = %transport.endpoint.id,
key_id = %transport.key.id,
"antigravity project metadata hydration response missing project"
);
return None;
};
let upstream_metadata = serde_json::json!({
"antigravity": {
"project_id": project_id,
"updated_at": current_unix_secs(),
}
});
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 !matches!(
provider_transport::antigravity::resolve_local_antigravity_request_auth(&hydrated),
provider_transport::antigravity::AntigravityRequestAuthSupport::Supported(_)
) {
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,
"antigravity project metadata hydration could not persist metadata"
);
}
Some(hydrated)
}
pub(crate) async fn hydrate_gemini_cli_project_metadata_for_transport(
&self,
transport: &GatewayProviderTransportSnapshot,
@@ -89,6 +193,30 @@ impl AppState {
}
}
fn extract_antigravity_load_code_assist_project_id(value: &Value) -> Option<String> {
let raw = value
.get("cloudaicompanionProject")
.or_else(|| value.get("cloudAiCompanionProject"))?;
if let Some(project_id) = raw
.as_str()
.map(str::trim)
.filter(|value| !value.is_empty())
{
return Some(project_id.to_string());
}
raw.as_object()
.and_then(|object| {
object
.get("id")
.or_else(|| object.get("project_id"))
.or_else(|| object.get("projectId"))
})
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
#[async_trait]
impl provider_transport::TransportTunnelAffinityLookup for AppState {
async fn lookup_tunnel_attachment_owner(