mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 16:37:46 +08:00
feat(codex): support agent identity accounts
This commit is contained in:
@@ -53,6 +53,130 @@ fn sanitize_windsurf_batch_import_error(error: &OAuthError) -> String {
|
||||
}
|
||||
}
|
||||
|
||||
fn copy_codex_agent_identity_field(
|
||||
auth_config: &mut Map<String, Value>,
|
||||
nested: &Map<String, Value>,
|
||||
canonical_key: &str,
|
||||
aliases: &[&str],
|
||||
) {
|
||||
if auth_config.contains_key(canonical_key) {
|
||||
return;
|
||||
}
|
||||
if let Some(value) = aliases.iter().find_map(|key| nested.get(*key)).cloned() {
|
||||
auth_config.insert(canonical_key.to_string(), value);
|
||||
}
|
||||
}
|
||||
|
||||
fn remove_codex_agent_identity_oauth_tokens(auth_config: &mut Map<String, Value>) {
|
||||
for key in [
|
||||
"access_token",
|
||||
"accessToken",
|
||||
"refresh_token",
|
||||
"refreshToken",
|
||||
"id_token",
|
||||
"idToken",
|
||||
"expires_at",
|
||||
"expiresAt",
|
||||
"expires_in",
|
||||
"expiresIn",
|
||||
] {
|
||||
auth_config.remove(key);
|
||||
}
|
||||
}
|
||||
|
||||
fn codex_agent_identity_auth_config_from_import(
|
||||
entry: &AdminProviderOAuthBatchImportEntry,
|
||||
) -> Result<Option<Map<String, Value>>, String> {
|
||||
let Some(raw_credentials) = entry.raw_credentials.as_ref() else {
|
||||
return Ok(None);
|
||||
};
|
||||
if !aether_provider_transport::is_codex_agent_identity_auth_config_value(raw_credentials) {
|
||||
return Ok(None);
|
||||
}
|
||||
let mut auth_config = raw_credentials
|
||||
.as_object()
|
||||
.cloned()
|
||||
.ok_or_else(|| "Agent Identity 凭据必须是 JSON 对象".to_string())?;
|
||||
remove_codex_agent_identity_oauth_tokens(&mut auth_config);
|
||||
for nested_key in ["agent_identity", "agentIdentity"] {
|
||||
if let Some(nested) = auth_config
|
||||
.get_mut(nested_key)
|
||||
.and_then(Value::as_object_mut)
|
||||
{
|
||||
remove_codex_agent_identity_oauth_tokens(nested);
|
||||
}
|
||||
}
|
||||
let nested = auth_config
|
||||
.get("agent_identity")
|
||||
.or_else(|| auth_config.get("agentIdentity"))
|
||||
.and_then(Value::as_object)
|
||||
.cloned();
|
||||
let root = auth_config.clone();
|
||||
for (canonical_key, aliases) in [
|
||||
(
|
||||
"agent_runtime_id",
|
||||
&["agent_runtime_id", "agentRuntimeId"][..],
|
||||
),
|
||||
(
|
||||
"agent_private_key",
|
||||
&["agent_private_key", "agentPrivateKey"][..],
|
||||
),
|
||||
("task_id", &["task_id", "taskId"][..]),
|
||||
(
|
||||
"account_id",
|
||||
&[
|
||||
"account_id",
|
||||
"accountId",
|
||||
"chatgpt_account_id",
|
||||
"chatgptAccountId",
|
||||
][..],
|
||||
),
|
||||
(
|
||||
"account_user_id",
|
||||
&[
|
||||
"account_user_id",
|
||||
"accountUserId",
|
||||
"chatgpt_account_user_id",
|
||||
"chatgptAccountUserId",
|
||||
][..],
|
||||
),
|
||||
(
|
||||
"user_id",
|
||||
&["user_id", "userId", "chatgpt_user_id", "chatgptUserId"][..],
|
||||
),
|
||||
("email", &["email"][..]),
|
||||
(
|
||||
"plan_type",
|
||||
&[
|
||||
"plan_type",
|
||||
"planType",
|
||||
"chatgpt_plan_type",
|
||||
"chatgptPlanType",
|
||||
][..],
|
||||
),
|
||||
("account_name", &["account_name", "accountName"][..]),
|
||||
(
|
||||
"is_fedramp",
|
||||
&[
|
||||
"is_fedramp",
|
||||
"chatgpt_account_is_fedramp",
|
||||
"chatgptAccountIsFedramp",
|
||||
][..],
|
||||
),
|
||||
] {
|
||||
if let Some(nested) = nested.as_ref() {
|
||||
copy_codex_agent_identity_field(&mut auth_config, nested, canonical_key, aliases);
|
||||
}
|
||||
copy_codex_agent_identity_field(&mut auth_config, &root, canonical_key, aliases);
|
||||
}
|
||||
auth_config.insert("provider_type".to_string(), json!("codex"));
|
||||
auth_config.insert("auth_mode".to_string(), json!("agentIdentity"));
|
||||
aether_provider_transport::validate_codex_agent_identity_auth_config(&Value::Object(
|
||||
auth_config.clone(),
|
||||
))?;
|
||||
Ok(Some(auth_config))
|
||||
}
|
||||
|
||||
pub(super) fn estimate_admin_provider_oauth_batch_import_total(
|
||||
provider_type: &str,
|
||||
raw_credentials: &str,
|
||||
@@ -103,6 +227,18 @@ async fn resolve_admin_provider_oauth_batch_import_tokens(
|
||||
entry: &AdminProviderOAuthBatchImportEntry,
|
||||
request_proxy: Option<ProxySnapshot>,
|
||||
) -> Result<AdminProviderOAuthResolvedBatchImport, String> {
|
||||
if provider_type.eq_ignore_ascii_case("codex") {
|
||||
if let Some(auth_config) = codex_agent_identity_auth_config_from_import(entry)? {
|
||||
return Ok(AdminProviderOAuthResolvedBatchImport {
|
||||
// Agent Identity signs an assertion for every request. The existing OAuth
|
||||
// record keeps a placeholder in its encrypted token column only.
|
||||
access_token: "__placeholder__".to_string(),
|
||||
auth_config,
|
||||
expires_at: None,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
let refresh_token = entry
|
||||
.refresh_token
|
||||
.as_deref()
|
||||
@@ -546,8 +682,12 @@ pub(super) async fn execute_admin_provider_oauth_batch_import(
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::sanitize_windsurf_batch_import_error;
|
||||
use super::super::parse::parse_admin_provider_oauth_batch_import_entries;
|
||||
use super::{
|
||||
codex_agent_identity_auth_config_from_import, sanitize_windsurf_batch_import_error,
|
||||
};
|
||||
use aether_oauth::core::OAuthError;
|
||||
use serde_json::json;
|
||||
|
||||
#[test]
|
||||
fn windsurf_batch_import_error_redacts_http_body() {
|
||||
@@ -572,4 +712,54 @@ mod tests {
|
||||
assert!(!detail.contains("sk-secret"));
|
||||
assert!(!detail.contains("secret-token"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalizes_codex_agent_identity_import_without_access_token() {
|
||||
let entries = parse_admin_provider_oauth_batch_import_entries(
|
||||
"codex",
|
||||
r#"{
|
||||
"type":"sub2api-data",
|
||||
"version":1,
|
||||
"accounts":[{
|
||||
"name":"[email protected]",
|
||||
"platform":"openai",
|
||||
"credentials":{
|
||||
"auth_mode":"agentIdentity",
|
||||
"id_token":"stale-id-token",
|
||||
"agent_identity":{
|
||||
"agent_runtime_id":"runtime-1",
|
||||
"agent_private_key":"MC4CAQAwBQYDK2VwBCIEIAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
|
||||
"accountId":"account-1",
|
||||
"chatgptUserId":"user-1",
|
||||
"chatgptAccountIsFedramp":true,
|
||||
"access_token":"stale-access-token"
|
||||
}
|
||||
}
|
||||
}]
|
||||
}"#,
|
||||
);
|
||||
|
||||
let auth_config = codex_agent_identity_auth_config_from_import(&entries[0])
|
||||
.expect("Agent Identity import should validate")
|
||||
.expect("Agent Identity config should be recognized");
|
||||
|
||||
assert_eq!(auth_config.get("provider_type"), Some(&json!("codex")));
|
||||
assert_eq!(auth_config.get("auth_mode"), Some(&json!("agentIdentity")));
|
||||
assert_eq!(
|
||||
auth_config.get("agent_runtime_id"),
|
||||
Some(&json!("runtime-1"))
|
||||
);
|
||||
assert_eq!(auth_config.get("account_id"), Some(&json!("account-1")));
|
||||
assert_eq!(auth_config.get("user_id"), Some(&json!("user-1")));
|
||||
assert_eq!(auth_config.get("is_fedramp"), Some(&json!(true)));
|
||||
assert_eq!(
|
||||
auth_config.get("account_name"),
|
||||
Some(&json!("[email protected]"))
|
||||
);
|
||||
assert!(!auth_config.contains_key("id_token"));
|
||||
assert!(auth_config
|
||||
.get("agent_identity")
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.is_some_and(|nested| !nested.contains_key("access_token")));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -215,6 +215,34 @@ fn extract_admin_provider_oauth_batch_import_entry(
|
||||
serde_json::Value::Object(object) => {
|
||||
let is_grok = provider_type.trim().eq_ignore_ascii_case("grok");
|
||||
let is_windsurf = provider_type.trim().eq_ignore_ascii_case("windsurf");
|
||||
let is_codex_agent_identity = provider_type.trim().eq_ignore_ascii_case("codex")
|
||||
&& aether_provider_transport::is_codex_agent_identity_auth_config_value(item);
|
||||
if is_codex_agent_identity {
|
||||
return Some(AdminProviderOAuthBatchImportEntry {
|
||||
parse_error: None,
|
||||
refresh_token: None,
|
||||
access_token: None,
|
||||
export_access_token: None,
|
||||
raw_credentials: Some(item.clone()),
|
||||
expires_at: None,
|
||||
account_id: None,
|
||||
account_user_id: None,
|
||||
plan_type: None,
|
||||
pool_tier: None,
|
||||
user_id: None,
|
||||
email: None,
|
||||
account_name: None,
|
||||
project_id: None,
|
||||
client_version: None,
|
||||
session_id: None,
|
||||
sso_rw_token: None,
|
||||
cf_cookies: None,
|
||||
cf_clearance: None,
|
||||
request_headers: None,
|
||||
user_agent: None,
|
||||
browser_profile: None,
|
||||
});
|
||||
}
|
||||
let refresh_token = coerce_admin_provider_oauth_import_str(
|
||||
object
|
||||
.get("refresh_token")
|
||||
@@ -434,6 +462,99 @@ fn extract_admin_provider_oauth_batch_import_entry(
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_sub2api_export_accounts(
|
||||
provider_type: &str,
|
||||
object: &serde_json::Map<String, serde_json::Value>,
|
||||
) -> Option<Vec<AdminProviderOAuthBatchImportEntry>> {
|
||||
if !provider_type.trim().eq_ignore_ascii_case("codex") {
|
||||
return None;
|
||||
}
|
||||
let is_sub2api_export = object
|
||||
.get("type")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.is_some_and(|value| value.trim().eq_ignore_ascii_case("sub2api-data"));
|
||||
if !is_sub2api_export {
|
||||
return None;
|
||||
}
|
||||
|
||||
let Some(accounts) = object.get("accounts").and_then(serde_json::Value::as_array) else {
|
||||
return Some(vec![parse_error_entry(
|
||||
"sub2api 导出缺少 accounts 数组".to_string(),
|
||||
)]);
|
||||
};
|
||||
|
||||
let mut entries = Vec::new();
|
||||
for (index, account) in accounts.iter().enumerate() {
|
||||
let Some(account) = account.as_object() else {
|
||||
entries.push(parse_error_entry(format!(
|
||||
"sub2api 第 {} 个账号必须是 JSON 对象",
|
||||
index + 1
|
||||
)));
|
||||
continue;
|
||||
};
|
||||
if account
|
||||
.get("platform")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.is_some_and(|platform| !platform.trim().eq_ignore_ascii_case("openai"))
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(mut credentials) = account
|
||||
.get("credentials")
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.cloned()
|
||||
else {
|
||||
entries.push(parse_error_entry(format!(
|
||||
"sub2api 第 {} 个账号缺少 credentials 对象",
|
||||
index + 1
|
||||
)));
|
||||
continue;
|
||||
};
|
||||
if let Some(name) = account
|
||||
.get("name")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
credentials
|
||||
.entry("account_name".to_string())
|
||||
.or_insert_with(|| json!(name));
|
||||
}
|
||||
if let Some(extra) = account.get("extra").and_then(serde_json::Value::as_object) {
|
||||
for key in [
|
||||
"account_id",
|
||||
"chatgpt_account_id",
|
||||
"chatgpt_user_id",
|
||||
"chatgpt_account_is_fedramp",
|
||||
"email",
|
||||
"plan_type",
|
||||
"workspace_id",
|
||||
] {
|
||||
if let Some(value) = extra.get(key).cloned() {
|
||||
credentials.entry(key.to_string()).or_insert(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let credentials = serde_json::Value::Object(credentials);
|
||||
match extract_admin_provider_oauth_batch_import_entry(provider_type, &credentials) {
|
||||
Some(entry) => entries.push(entry),
|
||||
None => entries.push(parse_error_entry(format!(
|
||||
"sub2api 第 {} 个账号没有可导入的凭据",
|
||||
index + 1
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
if entries.is_empty() {
|
||||
entries.push(parse_error_entry(
|
||||
"sub2api 导出中没有可导入的 OpenAI 账号".to_string(),
|
||||
));
|
||||
}
|
||||
Some(entries)
|
||||
}
|
||||
|
||||
pub(super) fn parse_admin_provider_oauth_batch_import_entries(
|
||||
provider_type: &str,
|
||||
raw_credentials: &str,
|
||||
@@ -459,12 +580,15 @@ pub(super) fn parse_admin_provider_oauth_batch_import_entries(
|
||||
}
|
||||
|
||||
if raw.starts_with('{') {
|
||||
if let Ok(value @ serde_json::Value::Object(_)) =
|
||||
serde_json::from_str::<serde_json::Value>(raw)
|
||||
{
|
||||
return extract_admin_provider_oauth_batch_import_entry(provider_type, &value)
|
||||
.into_iter()
|
||||
.collect();
|
||||
if let Ok(value) = serde_json::from_str::<serde_json::Value>(raw) {
|
||||
if let Some(object) = value.as_object() {
|
||||
if let Some(entries) = parse_sub2api_export_accounts(provider_type, object) {
|
||||
return entries;
|
||||
}
|
||||
return extract_admin_provider_oauth_batch_import_entry(provider_type, &value)
|
||||
.into_iter()
|
||||
.collect();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -766,6 +890,76 @@ mod tests {
|
||||
assert_eq!(entries[0].email.as_deref(), Some("[email protected]"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn preserves_codex_agent_identity_entry_without_access_token() {
|
||||
let entries = parse_admin_provider_oauth_batch_import_entries(
|
||||
"codex",
|
||||
r#"{"auth_mode":"agentIdentity","agent_identity":{"agent_runtime_id":"runtime-1","agent_private_key":"not-validated-until-import"}}"#,
|
||||
);
|
||||
|
||||
assert_eq!(entries.len(), 1);
|
||||
assert!(entries[0].refresh_token.is_none());
|
||||
assert!(entries[0].access_token.is_none());
|
||||
assert_eq!(
|
||||
entries[0]
|
||||
.raw_credentials
|
||||
.as_ref()
|
||||
.and_then(|value| value.get("auth_mode")),
|
||||
Some(&json!("agentIdentity"))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unwraps_sub2api_agent_identity_export_accounts() {
|
||||
let entries = parse_admin_provider_oauth_batch_import_entries(
|
||||
"codex",
|
||||
r#"{
|
||||
"type":"sub2api-data",
|
||||
"version":1,
|
||||
"accounts":[
|
||||
{
|
||||
"name":"[email protected]",
|
||||
"platform":"openai",
|
||||
"type":"oauth",
|
||||
"credentials":{
|
||||
"auth_mode":"agentIdentity",
|
||||
"agent_runtime_id":"runtime-1",
|
||||
"agent_private_key":"test-key",
|
||||
"task_id":"task-1",
|
||||
"chatgpt_account_id":"account-1"
|
||||
},
|
||||
"extra":{
|
||||
"email":"[email protected]",
|
||||
"chatgpt_user_id":"user-1"
|
||||
}
|
||||
},
|
||||
{
|
||||
"name":"[email protected]",
|
||||
"platform":"anthropic",
|
||||
"credentials":{"access_token":"ignored-token"}
|
||||
}
|
||||
]
|
||||
}"#,
|
||||
);
|
||||
|
||||
assert_eq!(entries.len(), 1);
|
||||
assert!(entries[0].parse_error.is_none());
|
||||
assert!(entries[0].refresh_token.is_none());
|
||||
assert!(entries[0].access_token.is_none());
|
||||
let credentials = entries[0]
|
||||
.raw_credentials
|
||||
.as_ref()
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.expect("Agent Identity credentials should be preserved");
|
||||
assert_eq!(credentials.get("auth_mode"), Some(&json!("agentIdentity")));
|
||||
assert_eq!(
|
||||
credentials.get("account_name"),
|
||||
Some(&json!("[email protected]"))
|
||||
);
|
||||
assert_eq!(credentials.get("email"), Some(&json!("[email protected]")));
|
||||
assert_eq!(credentials.get("chatgpt_user_id"), Some(&json!("user-1")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parses_codex_import_header_overrides() {
|
||||
let entries = parse_admin_provider_oauth_batch_import_entries(
|
||||
|
||||
@@ -338,7 +338,69 @@ pub(super) async fn execute_provider_quota_plan(
|
||||
quota_kind: &str,
|
||||
) -> Result<ProviderQuotaExecutionOutcome, GatewayError> {
|
||||
match state.execute_execution_runtime_sync_plan(None, &plan).await {
|
||||
Ok(result) => Ok(ProviderQuotaExecutionOutcome::Response(result)),
|
||||
Ok(result) => {
|
||||
if !crate::provider_transport::is_codex_agent_identity_transport(transport)
|
||||
|| !crate::provider_transport::is_codex_agent_identity_invalid_task_response(
|
||||
result.status_code,
|
||||
extract_execution_error_message(&result).as_deref(),
|
||||
)
|
||||
{
|
||||
return Ok(ProviderQuotaExecutionOutcome::Response(result));
|
||||
}
|
||||
|
||||
let refreshed_entry = match state.force_local_oauth_refresh_entry(transport).await {
|
||||
Ok(Some(entry)) => entry,
|
||||
Ok(None) => {
|
||||
return Ok(ProviderQuotaExecutionOutcome::Failure(
|
||||
"Agent Identity 任务重注册未返回认证信息".to_string(),
|
||||
));
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(
|
||||
key_id = %transport.key.id,
|
||||
endpoint_id = %transport.endpoint.id,
|
||||
quota_kind = %quota_kind,
|
||||
error = %error,
|
||||
"gateway Agent Identity quota task recovery failed"
|
||||
);
|
||||
return Ok(ProviderQuotaExecutionOutcome::Failure(format!(
|
||||
"Agent Identity 任务重注册失败: {error}"
|
||||
)));
|
||||
}
|
||||
};
|
||||
let header_name = refreshed_entry.auth_header_name.trim().to_ascii_lowercase();
|
||||
let header_value = refreshed_entry.auth_header_value.trim();
|
||||
if header_name.is_empty() || header_value.is_empty() {
|
||||
return Ok(ProviderQuotaExecutionOutcome::Failure(
|
||||
"Agent Identity 任务重注册未返回有效认证信息".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let mut retry_plan = plan.clone();
|
||||
retry_plan
|
||||
.headers
|
||||
.retain(|name, _| !name.eq_ignore_ascii_case(&header_name));
|
||||
retry_plan
|
||||
.headers
|
||||
.insert(header_name, header_value.to_string());
|
||||
match state
|
||||
.execute_execution_runtime_sync_plan(None, &retry_plan)
|
||||
.await
|
||||
{
|
||||
Ok(result) => Ok(ProviderQuotaExecutionOutcome::Response(result)),
|
||||
Err(error) => {
|
||||
let error = error.into_message();
|
||||
warn!(
|
||||
key_id = %transport.key.id,
|
||||
endpoint_id = %transport.endpoint.id,
|
||||
quota_kind = %quota_kind,
|
||||
error = %error,
|
||||
"gateway Agent Identity quota task recovery retry failed"
|
||||
);
|
||||
Ok(ProviderQuotaExecutionOutcome::Failure(error))
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
let error = err.into_message();
|
||||
let proxy_node_id = plan
|
||||
|
||||
@@ -52,6 +52,19 @@ pub(crate) async fn build_admin_create_provider_key_record(
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.cloned();
|
||||
|
||||
if auth_type == "oauth"
|
||||
&& provider.provider_type.trim().eq_ignore_ascii_case("codex")
|
||||
&& auth_config
|
||||
.as_ref()
|
||||
.is_some_and(aether_provider_transport::is_codex_agent_identity_auth_config_value)
|
||||
{
|
||||
aether_provider_transport::validate_codex_agent_identity_auth_config(
|
||||
auth_config
|
||||
.as_ref()
|
||||
.expect("Agent Identity auth_config was checked"),
|
||||
)?;
|
||||
}
|
||||
|
||||
match auth_type.as_str() {
|
||||
"service_account" if auth_config_object.is_none() => {
|
||||
return Err("Service Account 认证模式下 auth_config 为必填字段".to_string());
|
||||
|
||||
@@ -78,6 +78,19 @@ pub(crate) fn build_admin_update_provider_key_record_with_existing_keys(
|
||||
.and_then(serde_json::Value::as_object)
|
||||
.cloned();
|
||||
|
||||
if target_auth_type == "oauth"
|
||||
&& provider.provider_type.trim().eq_ignore_ascii_case("codex")
|
||||
&& auth_config
|
||||
.as_ref()
|
||||
.is_some_and(aether_provider_transport::is_codex_agent_identity_auth_config_value)
|
||||
{
|
||||
aether_provider_transport::validate_codex_agent_identity_auth_config(
|
||||
auth_config
|
||||
.as_ref()
|
||||
.expect("Agent Identity auth_config was checked"),
|
||||
)?;
|
||||
}
|
||||
|
||||
match target_auth_type.as_str() {
|
||||
"api_key" | "bearer" => {
|
||||
if let Some(api_key) = api_key_value
|
||||
|
||||
Reference in New Issue
Block a user