mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-09 04:30:20 +08:00
fix oauth identity and codex reset credits
This commit is contained in:
@@ -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,
|
||||
@@ -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}])
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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": "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(
|
||||
|
||||
@@ -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<String, serde_json::Value>,
|
||||
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!({
|
||||
|
||||
@@ -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<String>) -> 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<String>) -> 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<ProviderOAuthTokenSet, OAuthError> {
|
||||
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::<Value>(&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<crate::provider::ProviderOAuthTokenSet, crate::core::OAuthError> {
|
||||
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<Vec<OAuthHttpRequest>>,
|
||||
}
|
||||
|
||||
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<OAuthHttpResponse, crate::core::OAuthError> {
|
||||
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();
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<ReturnType<typeof
|
||||
}
|
||||
|
||||
function getCodexResetCredits(key: PoolKeyDetail) {
|
||||
return getQuotaSnapshotProviderType(key) === 'codex'
|
||||
? key.status_snapshot?.quota?.reset_credits ?? null
|
||||
if (getQuotaSnapshotProviderType(key) !== 'codex') return null
|
||||
|
||||
const snapshot = key.status_snapshot?.quota?.reset_credits
|
||||
const snapshotUpdatedAt = key.status_snapshot?.quota?.updated_at
|
||||
const snapshotDisplay = snapshot
|
||||
? {
|
||||
...(typeof snapshotUpdatedAt === 'number' ? { updated_at: snapshotUpdatedAt } : {}),
|
||||
reset_credits: snapshot,
|
||||
}
|
||||
: null
|
||||
return mergeCodexQuotaDisplays(snapshotDisplay, key.upstream_metadata?.codex)?.reset_credits ?? null
|
||||
}
|
||||
|
||||
function getCodexCredentialGeneration(key: PoolKeyDetail): string | null | undefined {
|
||||
|
||||
Reference in New Issue
Block a user