From fe8ff268df3a37eff9c2ab381f016427c1839470 Mon Sep 17 00:00:00 2001 From: ZheFox <77232781+zhefox@users.noreply.github.com> Date: Fri, 4 Sep 2026 11:34:41 +0800 Subject: [PATCH 1/3] fix oauth identity and codex reset credits --- .../admin/provider/oauth/quota/shared.rs | 66 ++++++- .../admin/provider/oauth/state/exchange.rs | 22 ++- .../src/tests/control/admin/oauth.rs | 169 ++++++++++++++++ crates/aether-admin/src/provider/quota.rs | 101 +++++++++- .../src/provider/providers/antigravity.rs | 185 +++++++++++++++++- .../src/provider/providers/mod.rs | 2 +- frontend/src/views/admin/PoolManagement.vue | 13 +- 7 files changed, 545 insertions(+), 13 deletions(-) diff --git a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs index 55f3d0837..648fc7803 100644 --- a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs +++ b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs @@ -570,6 +570,50 @@ pub(crate) async fn reserve_codex_account_reset( Ok(None) } +fn record_locally_consumed_codex_reset_credit( + codex: &mut serde_json::Map, + observed_at_unix_secs: u64, +) { + let Some(reset_credits) = codex + .get_mut("reset_credits") + .and_then(serde_json::Value::as_object_mut) + else { + return; + }; + let Some(available_count) = reset_credits + .get("available_count") + .and_then(admin_provider_quota_pure::coerce_json_u64) + else { + return; + }; + + reset_credits.insert( + "available_count".to_string(), + serde_json::json!(available_count.saturating_sub(1)), + ); + reset_credits.insert( + "updated_at".to_string(), + serde_json::json!(observed_at_unix_secs), + ); + reset_credits.insert( + "detail_source".to_string(), + serde_json::json!("local_consume"), + ); + reset_credits.insert( + "detail_status".to_string(), + serde_json::json!("pending_refresh"), + ); + reset_credits.remove("detail_error"); + if let Some(credits) = reset_credits + .get_mut("credits") + .and_then(serde_json::Value::as_array_mut) + { + if !credits.is_empty() { + credits.remove(0); + } + } +} + pub(crate) async fn complete_codex_account_reset( state: &AdminAppState<'_>, key_id: &str, @@ -635,6 +679,9 @@ pub(crate) async fn complete_codex_account_reset( generation: reservation.generation, outcome: outcome.to_string(), }; + if outcome == "reset" { + record_locally_consumed_codex_reset_credit(&mut codex, fence_unix_ms / 1_000); + } codex_reset_write_bounded_history(&mut codex, &terminal); if codex_reset_reservation_from_object(&codex).as_ref() == Some(reservation) { codex.remove(admin_provider_quota_pure::CODEX_QUOTA_ACCOUNT_RESET_RESERVATION_KEY); @@ -1671,7 +1718,19 @@ mod tests { .expect("key should build"); key.encrypted_auth_config = Some("auth-v1".to_string()); key.upstream_metadata = Some(json!({ - "codex": {"credential_generation": "credential-v1"} + "codex": { + "credential_generation": "credential-v1", + "reset_credits": { + "available_count": 2, + "updated_at": 100u64, + "detail_source": "wham_readonly", + "detail_status": "available", + "credits": [ + {"id": "credit-1", "expires_at": 20_000u64}, + {"id": "credit-2", "expires_at": 30_000u64} + ] + } + } })); let credential = ProviderCatalogKeyOAuthCredentialFence { encrypted_api_key: None, @@ -1926,6 +1985,11 @@ mod tests { codex["account_quota_reset_history"][0]["outcome"], json!("reset") ); + assert_eq!(codex["reset_credits"]["available_count"], json!(1u64)); + assert_eq!( + codex["reset_credits"]["credits"], + json!([{"id": "credit-2", "expires_at": 30_000u64}]) + ); } } diff --git a/apps/aether-gateway/src/handlers/admin/provider/oauth/state/exchange.rs b/apps/aether-gateway/src/handlers/admin/provider/oauth/state/exchange.rs index bfda465fb..b737510f9 100644 --- a/apps/aether-gateway/src/handlers/admin/provider/oauth/state/exchange.rs +++ b/apps/aether-gateway/src/handlers/admin/provider/oauth/state/exchange.rs @@ -4,8 +4,9 @@ use super::super::errors::{ use crate::handlers::admin::request::{AdminAppState, AdminProviderOAuthTemplate}; use aether_contracts::ProxySnapshot; use aether_oauth::provider::providers::{ - ClaudeCodeProviderOAuthAdapter, GenericProviderOAuthAdapter, CLAUDE_CODE_PROVIDER_TYPE, - CLAUDE_CODE_TOKEN_URL, CLAUDE_CODE_WEB_BASE_URL, + AntigravityProviderOAuthAdapter, ClaudeCodeProviderOAuthAdapter, GenericProviderOAuthAdapter, + ANTIGRAVITY_USER_INFO_URL, CLAUDE_CODE_PROVIDER_TYPE, CLAUDE_CODE_TOKEN_URL, + CLAUDE_CODE_WEB_BASE_URL, }; use aether_oauth::provider::{ ProviderOAuthCookieAuthorizationInput, ProviderOAuthService, ProviderOAuthTransportContext, @@ -43,7 +44,14 @@ fn provider_oauth_exchange_context( fn provider_oauth_service_for_template( template: AdminProviderOAuthTemplate, token_url: String, + antigravity_user_info_url: String, ) -> Result> { + if template.provider_type.eq_ignore_ascii_case("antigravity") { + let adapter = AntigravityProviderOAuthAdapter::default() + .with_token_url_override(token_url) + .with_user_info_url_override(antigravity_user_info_url); + return Ok(ProviderOAuthService::new().with_adapter(Arc::new(adapter))); + } GenericProviderOAuthAdapter::for_provider_type(template.provider_type) .map(|adapter| adapter.with_token_url_override(token_url)) .map(|adapter| ProviderOAuthService::new().with_adapter(Arc::new(adapter))) @@ -75,7 +83,10 @@ pub(crate) async fn exchange_admin_provider_oauth_code( proxy: Option, ) -> Result> { let token_url = state.provider_oauth_token_url(template.provider_type, template.token_url); - let service = provider_oauth_service_for_template(template, token_url)?; + let antigravity_user_info_url = + state.provider_oauth_token_url("antigravity_user_info", ANTIGRAVITY_USER_INFO_URL); + let service = + provider_oauth_service_for_template(template, token_url, antigravity_user_info_url)?; let ctx = provider_oauth_exchange_context(template.provider_type, proxy); let executor = crate::oauth::GatewayOAuthHttpExecutor::new(*state); let result = service @@ -103,7 +114,10 @@ pub(crate) async fn exchange_admin_provider_oauth_refresh_token( proxy: Option, ) -> Result> { let token_url = state.provider_oauth_token_url(template.provider_type, template.token_url); - let service = provider_oauth_service_for_template(template, token_url)?; + let antigravity_user_info_url = + state.provider_oauth_token_url("antigravity_user_info", ANTIGRAVITY_USER_INFO_URL); + let service = + provider_oauth_service_for_template(template, token_url, antigravity_user_info_url)?; let ctx = provider_oauth_exchange_context(template.provider_type, proxy); let executor = crate::oauth::GatewayOAuthHttpExecutor::new(*state); let input = aether_oauth::provider::ProviderOAuthImportInput { diff --git a/apps/aether-gateway/src/tests/control/admin/oauth.rs b/apps/aether-gateway/src/tests/control/admin/oauth.rs index 8f330ca32..864680d45 100644 --- a/apps/aether-gateway/src/tests/control/admin/oauth.rs +++ b/apps/aether-gateway/src/tests/control/admin/oauth.rs @@ -3918,6 +3918,175 @@ async fn gateway_completes_admin_provider_oauth_provider_locally_with_trusted_ad upstream_handle.abort(); } +#[test] +fn gateway_names_new_antigravity_oauth_account_from_google_userinfo_email() { + run_admin_oauth_test( + "gateway_names_new_antigravity_oauth_account_from_google_userinfo_email", + gateway_names_new_antigravity_oauth_account_from_google_userinfo_email_impl, + ); +} + +async fn gateway_names_new_antigravity_oauth_account_from_google_userinfo_email_impl() { + let upstream_hits = Arc::new(Mutex::new(0usize)); + let upstream_hits_clone = Arc::clone(&upstream_hits); + let upstream = Router::new().fallback(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 token_hits = Arc::new(Mutex::new(0usize)); + let token_hits_clone = Arc::clone(&token_hits); + let user_info_hits = Arc::new(Mutex::new(0usize)); + let user_info_hits_clone = Arc::clone(&user_info_hits); + let seen_user_info_authorization = Arc::new(Mutex::new(None::)); + let seen_user_info_authorization_clone = Arc::clone(&seen_user_info_authorization); + let google_server = Router::new() + .route( + "/oauth/token", + post(move || { + let token_hits_inner = Arc::clone(&token_hits_clone); + async move { + *token_hits_inner.lock().expect("mutex should lock") += 1; + Json(json!({ + "access_token": "antigravity-access-token", + "refresh_token": "antigravity-refresh-token", + "token_type": "Bearer", + "expires_in": 3600, + "scope": "https://www.googleapis.com/auth/userinfo.email" + })) + } + }), + ) + .route( + "/oauth/userinfo", + get(move |headers: HeaderMap| { + let user_info_hits_inner = Arc::clone(&user_info_hits_clone); + let seen_authorization_inner = Arc::clone(&seen_user_info_authorization_clone); + async move { + *user_info_hits_inner.lock().expect("mutex should lock") += 1; + *seen_authorization_inner.lock().expect("mutex should lock") = headers + .get(http::header::AUTHORIZATION) + .and_then(|value| value.to_str().ok()) + .map(ToOwned::to_owned); + Json(json!({ + "email": "new-antigravity@example.com", + "verified_email": true, + "name": "Antigravity User" + })) + } + }), + ); + + let mut provider = sample_provider("provider-antigravity", "antigravity", 10); + provider.provider_type = "antigravity".to_string(); + let endpoint = sample_endpoint( + "endpoint-antigravity", + "provider-antigravity", + "gemini:generate_content", + "https://daily-cloudcode-pa.googleapis.com", + ); + let provider_catalog_repository = Arc::new(InMemoryProviderCatalogReadRepository::seed( + vec![provider], + vec![endpoint], + vec![], + )); + + let (upstream_url, upstream_handle) = start_server(upstream).await; + let (google_url, google_handle) = start_server(google_server).await; + let gateway = build_router_with_state( + AppState::new() + .expect("gateway should build") + .with_data_state_for_tests( + GatewayDataState::with_provider_catalog_repository_for_tests( + provider_catalog_repository.clone(), + ) + .with_encryption_key_for_tests(DEVELOPMENT_ENCRYPTION_KEY), + ) + .with_provider_oauth_state_entry_for_tests( + "nonce-antigravity-123", + json!({ + "nonce": "nonce-antigravity-123", + "key_id": "", + "provider_id": "provider-antigravity", + "provider_type": "antigravity", + "pkce_verifier": "verifier-antigravity-123", + }), + ) + .with_provider_oauth_token_url_for_tests( + "antigravity", + format!("{google_url}/oauth/token"), + ) + .with_provider_oauth_token_url_for_tests( + "antigravity_user_info", + format!("{google_url}/oauth/userinfo"), + ), + ); + let (gateway_url, gateway_handle) = start_server(gateway).await; + + let response = reqwest::Client::new() + .post(format!( + "{gateway_url}/api/admin/provider-oauth/providers/provider-antigravity/complete" + )) + .header(crate::constants::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!({ + "callback_url": "http://localhost:51121/oauth2callback?code=antigravity-code-123&state=nonce-antigravity-123" + })) + .send() + .await + .expect("request should succeed"); + + let status = response.status(); + let payload: Value = response.json().await.expect("json body should parse"); + assert_eq!(status, StatusCode::OK, "payload={payload}"); + assert_eq!(payload["provider_type"], "antigravity"); + assert_eq!(payload["email"], "new-antigravity@example.com"); + assert_eq!(payload["replaced"], false); + assert_eq!(*token_hits.lock().expect("mutex should lock"), 1); + assert_eq!(*user_info_hits.lock().expect("mutex should lock"), 1); + assert_eq!( + seen_user_info_authorization + .lock() + .expect("mutex should lock") + .as_deref(), + Some("Bearer antigravity-access-token") + ); + assert_eq!(*upstream_hits.lock().expect("mutex should lock"), 0); + + let key_id = payload["key_id"] + .as_str() + .expect("created key id should be returned") + .to_string(); + let persisted_keys = provider_catalog_repository + .list_keys_by_ids(std::slice::from_ref(&key_id)) + .await + .expect("created key should load"); + let persisted = persisted_keys.first().expect("created key should exist"); + assert_eq!(persisted.name, "new-antigravity@example.com"); + let decrypted_auth_config = decrypt_python_fernet_ciphertext( + DEVELOPMENT_ENCRYPTION_KEY, + persisted + .encrypted_auth_config + .as_deref() + .expect("auth config should be stored"), + ) + .expect("auth config should decrypt"); + let auth_config: Value = + serde_json::from_str(&decrypted_auth_config).expect("auth config json should parse"); + assert_eq!(auth_config["email"], "new-antigravity@example.com"); + assert_eq!(auth_config["refresh_token"], "antigravity-refresh-token"); + + gateway_handle.abort(); + google_handle.abort(); + upstream_handle.abort(); + drop(upstream_url); +} + #[test] fn gateway_imports_admin_provider_oauth_refresh_token_locally_with_trusted_admin_principal() { run_admin_oauth_test( diff --git a/crates/aether-admin/src/provider/quota.rs b/crates/aether-admin/src/provider/quota.rs index 349f26855..4603f65af 100644 --- a/crates/aether-admin/src/provider/quota.rs +++ b/crates/aether-admin/src/provider/quota.rs @@ -1887,6 +1887,46 @@ fn codex_quota_is_account_status_key(key: &str) -> bool { matches!(key, "allowed" | "limit_reached") } +fn codex_quota_merge_reset_credits( + current_object: &serde_json::Map, + incoming: &serde_json::Value, +) -> serde_json::Value { + let Some(incoming_object) = incoming.as_object() else { + return incoming.clone(); + }; + let mut merged = current_object + .get("reset_credits") + .and_then(serde_json::Value::as_object) + .cloned() + .unwrap_or_default(); + let failed_detail = incoming_object + .get("detail_status") + .and_then(serde_json::Value::as_str) + .is_some_and(|status| status.trim().eq_ignore_ascii_case("failed")); + if !failed_detail { + return incoming.clone(); + } + + for (key, value) in incoming_object { + // A failed readonly-detail request contributes diagnostics, not an + // authoritative empty list. Keep the last known items and count so a + // transient 429 cannot make reset credits disappear from the UI. + if failed_detail + && key == "credits" + && value.as_array().is_some_and(|credits| credits.is_empty()) + && merged + .get("credits") + .and_then(serde_json::Value::as_array) + .is_some_and(|credits| !credits.is_empty()) + { + continue; + } + merged.insert(key.clone(), value.clone()); + } + + serde_json::Value::Object(merged) +} + /// Merge a parsed Codex quota observation into the stored flat metadata. /// /// Positive `window_minutes` values identify windows independently of the @@ -1961,7 +2001,12 @@ pub fn merge_codex_quota_metadata_snapshot( { continue; } - merged.insert(key.clone(), value.clone()); + let value = if key == "reset_credits" { + codex_quota_merge_reset_credits(¤t_object, value) + } else { + value.clone() + }; + merged.insert(key.clone(), value); } if let Some(incoming_order) = context.request_order().filter(|incoming| { codex_quota_request_order_is_newer(*incoming, stored_metadata_watermark) @@ -4777,6 +4822,60 @@ mod tests { assert_eq!(outcome.metadata["primary_reset_at"], json!(20_000u64)); } + #[test] + fn codex_quota_failed_reset_credit_detail_preserves_last_known_count_and_items() { + let current = json!({ + "reset_credits": { + "available_count": 2, + "updated_at": 100u64, + "detail_source": "wham_readonly", + "detail_status": "available", + "credits": [{ + "id": "credit-1", + "display_key": "credit", + "status": "available", + "expires_at": 20_000u64 + }] + }, + "updated_at": 100u64 + }); + let incoming = json!({ + "reset_credits": { + "updated_at": 110u64, + "detail_source": "wham_readonly", + "detail_status": "failed", + "detail_error": "HTTP 429", + "credits": [] + } + }); + + let outcome = merge_codex_quota( + Some(¤t), + &incoming, + 110, + 110_000, + CodexQuotaWindowCoverage::Patch, + ); + + assert!(outcome.changed); + assert_eq!( + outcome.metadata["reset_credits"]["available_count"], + json!(2u64) + ); + assert_eq!( + outcome.metadata["reset_credits"]["credits"][0]["id"], + json!("credit-1") + ); + assert_eq!( + outcome.metadata["reset_credits"]["detail_status"], + json!("failed") + ); + assert_eq!( + outcome.metadata["reset_credits"]["detail_error"], + json!("HTTP 429") + ); + } + #[test] fn codex_quota_explicit_reset_allows_usage_drop_with_same_deadline() { let current = json!({ diff --git a/crates/aether-oauth/src/provider/providers/antigravity.rs b/crates/aether-oauth/src/provider/providers/antigravity.rs index 9bc86c053..644a82812 100644 --- a/crates/aether-oauth/src/provider/providers/antigravity.rs +++ b/crates/aether-oauth/src/provider/providers/antigravity.rs @@ -1,11 +1,18 @@ use super::generic::{ provider_account_state_from_metadata, template_for_provider_type, GenericProviderOAuthAdapter, }; -use crate::provider::ProviderOAuthAdapter; +use crate::core::OAuthError; +use crate::network::{OAuthHttpExecutor, OAuthHttpRequest}; +use crate::provider::{ProviderOAuthAdapter, ProviderOAuthTokenSet, ProviderOAuthTransportContext}; +use serde_json::Value; +use std::collections::BTreeMap; + +pub const ANTIGRAVITY_USER_INFO_URL: &str = "https://www.googleapis.com/oauth2/v2/userinfo"; #[derive(Debug, Clone)] pub struct AntigravityProviderOAuthAdapter { inner: GenericProviderOAuthAdapter, + user_info_url: String, } impl Default for AntigravityProviderOAuthAdapter { @@ -15,10 +22,95 @@ impl Default for AntigravityProviderOAuthAdapter { template_for_provider_type("antigravity") .expect("antigravity template should exist"), ), + user_info_url: ANTIGRAVITY_USER_INFO_URL.to_string(), } } } +impl AntigravityProviderOAuthAdapter { + pub fn with_token_url_override(mut self, token_url: impl Into) -> Self { + self.inner = self.inner.with_token_url_override(token_url); + self + } + + pub fn with_user_info_url_override(mut self, user_info_url: impl Into) -> Self { + self.user_info_url = user_info_url.into(); + self + } + + async fn enrich_google_identity( + &self, + executor: &dyn OAuthHttpExecutor, + ctx: &ProviderOAuthTransportContext, + mut result: ProviderOAuthTokenSet, + ) -> Result { + if result + .auth_config + .get("email") + .and_then(Value::as_str) + .is_some_and(|email| !email.trim().is_empty()) + { + return Ok(result); + } + + let response = executor + .execute(OAuthHttpRequest { + request_id: "provider-oauth:antigravity-user-info".to_string(), + method: reqwest::Method::GET, + url: self.user_info_url.clone(), + headers: BTreeMap::from([ + ("accept".to_string(), "application/json".to_string()), + ( + "authorization".to_string(), + result.token_set.bearer_header_value(), + ), + ]), + content_type: None, + json_body: None, + body_bytes: None, + network: ctx.network.clone(), + transport_profile: None, + }) + .await?; + if !(200..300).contains(&response.status_code) { + return Err(OAuthError::HttpStatus { + status_code: response.status_code, + body_excerpt: response.body_text.trim().chars().take(500).collect(), + }); + } + + let profile = response + .json_body + .or_else(|| serde_json::from_str::(&response.body_text).ok()) + .ok_or_else(|| OAuthError::invalid_response("userinfo response is not json"))?; + if profile.get("verified_email").and_then(Value::as_bool) == Some(false) { + return Err(OAuthError::invalid_response( + "userinfo response returned an unverified email", + )); + } + let email = profile + .get("email") + .and_then(Value::as_str) + .map(str::trim) + .filter(|email| !email.is_empty()) + .ok_or_else(|| OAuthError::invalid_response("userinfo response missing email"))? + .to_string(); + + if let Some(auth_config) = result.auth_config.as_object_mut() { + auth_config.insert("email".to_string(), Value::String(email.clone())); + } + if let Some(token_payload) = result + .token_set + .raw_payload + .as_mut() + .and_then(Value::as_object_mut) + { + token_payload.insert("email".to_string(), Value::String(email)); + } + Ok(result) + } +} + #[async_trait::async_trait] impl ProviderOAuthAdapter for AntigravityProviderOAuthAdapter { fn provider_type(&self) -> &'static str { @@ -59,9 +151,11 @@ impl ProviderOAuthAdapter for AntigravityProviderOAuthAdapter { state: &str, pkce_verifier: Option<&str>, ) -> Result { - self.inner + let result = self + .inner .exchange_code(executor, ctx, code, state, pkce_verifier) - .await + .await?; + self.enrich_google_identity(executor, ctx, result).await } async fn import_credentials( @@ -111,7 +205,7 @@ impl ProviderOAuthAdapter for AntigravityProviderOAuthAdapter { #[cfg(test)] mod tests { - use super::AntigravityProviderOAuthAdapter; + use super::{AntigravityProviderOAuthAdapter, ANTIGRAVITY_USER_INFO_URL}; use crate::network::{OAuthHttpExecutor, OAuthHttpRequest, OAuthHttpResponse}; use crate::provider::{ ProviderOAuthAccount, ProviderOAuthAdapter, ProviderOAuthTransportContext, @@ -119,9 +213,15 @@ mod tests { use async_trait::async_trait; use serde_json::json; use std::collections::BTreeMap; + use std::sync::Mutex; struct UnusedExecutor; + #[derive(Default)] + struct GoogleOAuthExecutor { + requests: Mutex>, + } + fn transport_context() -> ProviderOAuthTransportContext { ProviderOAuthTransportContext { provider_id: String::new(), @@ -148,6 +248,43 @@ mod tests { } } + #[async_trait] + impl OAuthHttpExecutor for GoogleOAuthExecutor { + async fn execute( + &self, + request: OAuthHttpRequest, + ) -> Result { + let request_id = request.request_id.clone(); + self.requests + .lock() + .expect("requests should lock") + .push(request); + match request_id.as_str() { + "provider-oauth:exchange-code" => Ok(OAuthHttpResponse { + status_code: 200, + body_text: json!({ + "access_token": "google-access-token", + "refresh_token": "google-refresh-token", + "token_type": "Bearer", + "expires_in": 3600 + }) + .to_string(), + json_body: None, + }), + "provider-oauth:antigravity-user-info" => Ok(OAuthHttpResponse { + status_code: 200, + body_text: json!({ + "email": "antigravity@example.com", + "verified_email": true + }) + .to_string(), + json_body: None, + }), + other => panic!("unexpected OAuth request: {other}"), + } + } + } + #[test] fn antigravity_authorize_requests_offline_refresh_token() { let adapter = AntigravityProviderOAuthAdapter::default(); @@ -171,6 +308,46 @@ mod tests { ); } + #[tokio::test] + async fn antigravity_exchange_fetches_google_email_for_account_identity() { + let adapter = AntigravityProviderOAuthAdapter::default(); + let ctx = transport_context(); + let executor = GoogleOAuthExecutor::default(); + + let result = adapter + .exchange_code( + &executor, + &ctx, + "authorization-code", + "state-1", + Some("verifier-1"), + ) + .await + .expect("Antigravity OAuth exchange should succeed"); + + assert_eq!( + result.auth_config.get("email"), + Some(&json!("antigravity@example.com")) + ); + assert_eq!( + result + .token_set + .raw_payload + .as_ref() + .and_then(|payload| payload.get("email")), + Some(&json!("antigravity@example.com")) + ); + let requests = executor.requests.lock().expect("requests should lock"); + assert_eq!(requests.len(), 2); + assert_eq!(requests[1].url, ANTIGRAVITY_USER_INFO_URL); + assert_eq!(requests[1].method, reqwest::Method::GET); + assert_eq!( + requests[1].headers.get("authorization").map(String::as_str), + Some("Bearer google-access-token") + ); + assert_eq!(requests[1].network, ctx.network); + } + #[tokio::test] async fn antigravity_probe_marks_forbidden_metadata_invalid() { let adapter = AntigravityProviderOAuthAdapter::default(); diff --git a/crates/aether-oauth/src/provider/providers/mod.rs b/crates/aether-oauth/src/provider/providers/mod.rs index 4fe916a75..07f752f36 100644 --- a/crates/aether-oauth/src/provider/providers/mod.rs +++ b/crates/aether-oauth/src/provider/providers/mod.rs @@ -5,7 +5,7 @@ mod generic; mod kiro; mod windsurf; -pub use antigravity::AntigravityProviderOAuthAdapter; +pub use antigravity::{AntigravityProviderOAuthAdapter, ANTIGRAVITY_USER_INFO_URL}; pub use claude_code::{ ClaudeCodeProviderOAuthAdapter, CLAUDE_CODE_AUTHORIZE_URL, CLAUDE_CODE_CLIENT_ID, CLAUDE_CODE_COOKIE_SCOPE, CLAUDE_CODE_OAUTH_SCOPES, CLAUDE_CODE_PROVIDER_TYPE, diff --git a/frontend/src/views/admin/PoolManagement.vue b/frontend/src/views/admin/PoolManagement.vue index 684e7a7b8..843ab750f 100644 --- a/frontend/src/views/admin/PoolManagement.vue +++ b/frontend/src/views/admin/PoolManagement.vue @@ -1153,6 +1153,7 @@ import { getCodexResetCreditAvailableCount, getCodexResetCreditReservationIdempotencyKey, getVisibleCodexResetCreditItems, + mergeCodexQuotaDisplays, readPendingCodexResetCreditIdempotencyKey, rememberPendingCodexResetCreditIdempotencyKey, } from '@/features/providers/components/codex-reset-credit-display' @@ -2235,9 +2236,17 @@ function applyQuotaRefreshResultToCurrentPage(result: Awaited Date: Fri, 4 Sep 2026 12:06:39 +0800 Subject: [PATCH 2/3] fix(antigravity): sync discovered models into catalog --- .../admin/provider/oauth/quota/antigravity.rs | 62 +++++++++++++++++++ .../tests/control/admin/endpoints/quota.rs | 29 +++++++++ crates/aether-model-fetch/src/lib.rs | 5 +- crates/aether-model-fetch/src/strategy.rs | 10 ++- 4 files changed, 103 insertions(+), 3 deletions(-) diff --git a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/antigravity.rs b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/antigravity.rs index d3e26fc79..2cbd0848a 100644 --- a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/antigravity.rs +++ b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/antigravity.rs @@ -5,6 +5,7 @@ use super::shared::{ quota_key_auto_removed, quota_refresh_success_invalid_state, resolve_provider_quota_execution_timeouts, ProviderQuotaExecutionOutcome, }; +use crate::handlers::admin::provider::shared::payloads::AdminImportProviderModelsRequest; use crate::handlers::admin::request::{AdminAppState, AdminGatewayProviderTransportSnapshot}; use crate::GatewayError; use aether_admin::provider::quota::{ @@ -22,6 +23,63 @@ use std::collections::BTreeMap; use std::time::{SystemTime, UNIX_EPOCH}; use tracing::warn; +fn antigravity_discovered_model_ids(metadata_update: Option<&serde_json::Value>) -> Vec { + metadata_update + .and_then(|value| value.pointer("/antigravity/quota_by_model")) + .and_then(serde_json::Value::as_object) + .into_iter() + .flat_map(|models| models.keys()) + .map(String::as_str) + .filter(|model_id| aether_model_fetch::antigravity_model_id_is_routable(model_id)) + .map(ToOwned::to_owned) + .collect() +} + +async fn sync_antigravity_discovered_models( + state: &AdminAppState<'_>, + provider_id: &str, + metadata_update: Option<&serde_json::Value>, +) { + if !state.has_global_model_data_reader() || !state.has_global_model_data_writer() { + return; + } + let model_ids = antigravity_discovered_model_ids(metadata_update); + if model_ids.is_empty() { + return; + } + + let result = state + .build_admin_import_provider_models_payload( + provider_id, + AdminImportProviderModelsRequest { + model_ids, + tiered_pricing: None, + price_per_request: None, + }, + ) + .await; + match result { + Ok(payload) => { + let errors = payload + .get("errors") + .and_then(serde_json::Value::as_array) + .map(Vec::len) + .unwrap_or(0); + if errors > 0 { + warn!( + provider_id, + errors, "Antigravity discovered-model catalog sync completed with item errors" + ); + } + } + Err(error) => warn!( + provider_id, + error = %error, + "Antigravity discovered-model catalog sync failed" + ), + } +} + async fn execute_antigravity_quota_plan( state: &AdminAppState<'_>, transport: &AdminGatewayProviderTransportSnapshot, @@ -330,6 +388,10 @@ pub(crate) async fn refresh_antigravity_provider_quota_locally( continue; } + if status == "success" { + sync_antigravity_discovered_models(state, &provider.id, metadata_update.as_ref()).await; + } + if status == "success" { success_count += 1; } else { diff --git a/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs b/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs index 444d550f2..7120cf871 100644 --- a/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs +++ b/apps/aether-gateway/src/tests/control/admin/endpoints/quota.rs @@ -4,8 +4,12 @@ use std::sync::{Arc, Mutex}; use aether_crypto::{ decrypt_python_fernet_ciphertext, encrypt_python_fernet_plaintext, DEVELOPMENT_ENCRYPTION_KEY, }; +use aether_data::repository::global_models::InMemoryGlobalModelReadRepository; use aether_data::repository::provider_catalog::InMemoryProviderCatalogReadRepository; use aether_data::repository::proxy_nodes::InMemoryProxyNodeRepository; +use aether_data_contracts::repository::global_models::{ + AdminProviderModelListQuery, GlobalModelReadRepository, +}; use aether_data_contracts::repository::provider_catalog::{ ProviderCatalogReadRepository, StoredProviderCatalogKey, StoredProviderCatalogProvider, }; @@ -2330,6 +2334,12 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru }, "gemini-2.5-pro": { "displayName": "Gemini 2.5 Pro" + }, + "gemini-3.7-flash-tiered": { + "displayName": "Gemini 3.7 Flash" + }, + "chat_23310": { + "displayName": "Internal Chat" } } }), @@ -2431,6 +2441,7 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru )], vec![key], )); + let global_model_repository = Arc::new(InMemoryGlobalModelReadRepository::default()); let (upstream_url, upstream_handle) = start_server(upstream).await; let (execution_runtime_url, execution_runtime_handle) = start_server(execution_runtime).await; @@ -2440,6 +2451,7 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru GatewayDataState::with_provider_catalog_repository_for_tests( provider_catalog_repository.clone(), ) + .with_global_model_repository_for_tests(global_model_repository.clone()) .with_encryption_key_for_tests(DEVELOPMENT_ENCRYPTION_KEY), ), ); @@ -2541,6 +2553,23 @@ async fn gateway_refreshes_admin_provider_quota_locally_for_antigravity_with_tru .and_then(|value| value.get("remaining_fraction")), Some(&json!(0.25)) ); + let imported_provider_models = global_model_repository + .list_admin_provider_models(&AdminProviderModelListQuery { + provider_id: "provider-antigravity".to_string(), + is_active: None, + offset: 0, + limit: 100, + }) + .await + .expect("imported Antigravity provider models should read"); + let imported_model_names = imported_provider_models + .iter() + .map(|model| model.provider_model_name.as_str()) + .collect::>(); + assert!(imported_model_names.contains("claude-sonnet-4")); + assert!(imported_model_names.contains("gemini-2.5-pro")); + assert!(imported_model_names.contains("gemini-3.7-flash-tiered")); + assert!(!imported_model_names.contains("chat_23310")); assert_eq!( reloaded[0] .upstream_metadata diff --git a/crates/aether-model-fetch/src/lib.rs b/crates/aether-model-fetch/src/lib.rs index bb3432805..8ef4b5855 100644 --- a/crates/aether-model-fetch/src/lib.rs +++ b/crates/aether-model-fetch/src/lib.rs @@ -21,8 +21,9 @@ pub use logic::{ upstream_metadata_namespace_updates, ModelFetchRunSummary, ModelsFetchPage, ModelsFetchSuccess, }; pub use strategy::{ - fetch_models_from_transports, fetch_models_from_transports_for_client_version, - ModelFetchStrategy, ModelFetchStrategyKind, ModelsFetchOutcome, SelectedModelFetchStrategy, + antigravity_model_id_is_routable, fetch_models_from_transports, + fetch_models_from_transports_for_client_version, ModelFetchStrategy, ModelFetchStrategyKind, + ModelsFetchOutcome, SelectedModelFetchStrategy, }; pub use transport::{ build_antigravity_fetch_available_models_plan, build_antigravity_load_code_assist_plan, diff --git a/crates/aether-model-fetch/src/strategy.rs b/crates/aether-model-fetch/src/strategy.rs index 64951adb8..249bec490 100644 --- a/crates/aether-model-fetch/src/strategy.rs +++ b/crates/aether-model-fetch/src/strategy.rs @@ -1018,7 +1018,7 @@ fn parse_antigravity_models_response(body: &Value) -> Result<(Vec, Option let mut quota_by_model = serde_json::Map::new(); for (model_id, model_data) in models_object { let model_id = model_id.trim(); - if model_id.is_empty() || ANTIGRAVITY_BLOCKED_MODELS.contains(&model_id) { + if !antigravity_model_id_is_routable(model_id) { continue; } let model_object = model_data.as_object().cloned().unwrap_or_default(); @@ -1054,6 +1054,14 @@ fn parse_antigravity_models_response(body: &Value) -> Result<(Vec, Option Ok((models, upstream_metadata)) } +pub fn antigravity_model_id_is_routable(model_id: &str) -> bool { + let model_id = model_id.trim(); + !model_id.is_empty() + && !ANTIGRAVITY_BLOCKED_MODELS + .iter() + .any(|blocked| blocked.eq_ignore_ascii_case(model_id)) +} + fn parse_kiro_available_models_response( body: &Value, ) -> Result<(Vec, Option), String> { From c8d1ae3e7ea4fb8021f0bd3cf2c59f0484c055c6 Mon Sep 17 00:00:00 2001 From: ZheFox <77232781+zhefox@users.noreply.github.com> Date: Fri, 4 Sep 2026 13:02:22 +0800 Subject: [PATCH 3/3] test(codex): preserve reset credit fixture metadata --- .../handlers/admin/provider/oauth/quota/shared.rs | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs index 648fc7803..9baaa7d39 100644 --- a/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs +++ b/apps/aether-gateway/src/handlers/admin/provider/oauth/quota/shared.rs @@ -1824,6 +1824,13 @@ mod tests { let key_id = "key-codex-reset-credential-generation"; let (app, repository, credential) = codex_reset_state_machine_test_state(key_id); let admin_state = AdminAppState::new(&app); + let original_metadata = repository + .list_keys_by_ids(&[key_id.to_string()]) + .await + .expect("key should load before reservation") + .pop() + .expect("key should exist before reservation") + .upstream_metadata; let result = reserve_codex_account_reset( &admin_state, @@ -1847,10 +1854,7 @@ mod tests { .expect("key should reload") .pop() .expect("key should exist"); - assert_eq!( - stored.upstream_metadata.unwrap()["codex"], - json!({"credential_generation":"credential-v1"}) - ); + assert_eq!(stored.upstream_metadata, original_metadata); } #[tokio::test]