Merge pull request #794 from zhefox/codex/antigravity-import-email

Codex/antigravity import email
This commit is contained in:
ZheFox
2026-09-04 13:11:08 +08:00
committed by GitHub
11 changed files with 656 additions and 20 deletions
@@ -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<String> {
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 {
@@ -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<String, serde_json::Value>,
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,
@@ -1765,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,
@@ -1788,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]
@@ -1926,6 +1989,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}])
);
}
}
@@ -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<ProviderOAuthService, Response<Body>> {
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<ProxySnapshot>,
) -> Result<serde_json::Value, Response<Body>> {
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<ProxySnapshot>,
) -> Result<serde_json::Value, Response<Body>> {
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 {
@@ -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::<std::collections::BTreeSet<_>>();
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
@@ -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::<String>));
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": "[email protected]",
"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"], "[email protected]");
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, "[email protected]");
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"], "[email protected]");
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(