mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 16:37:46 +08:00
增加 Codex 重置次数功能
This commit is contained in:
@@ -195,6 +195,17 @@ pub(super) fn classify_admin_endpoints_family_route(
|
||||
"admin:endpoints_manage",
|
||||
false,
|
||||
))
|
||||
} else if method == http::Method::POST
|
||||
&& normalized_path.starts_with("/api/admin/endpoints/keys/")
|
||||
&& normalized_path.ends_with("/codex-reset-credit/consume")
|
||||
{
|
||||
Some(classified(
|
||||
"admin_proxy",
|
||||
"endpoints_manage",
|
||||
"codex_reset_credit_consume",
|
||||
"admin:endpoints_manage",
|
||||
false,
|
||||
))
|
||||
} else if method == http::Method::POST
|
||||
&& normalized_path.starts_with("/api/admin/endpoints/providers/")
|
||||
&& normalized_path.ends_with("/refresh-quota")
|
||||
|
||||
@@ -437,6 +437,23 @@ fn classifies_admin_refresh_provider_quota_as_admin_proxy_route() {
|
||||
assert!(!decision.is_execution_runtime_candidate());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classifies_admin_codex_reset_credit_consume_as_admin_proxy_route() {
|
||||
let headers = http::HeaderMap::new();
|
||||
let uri: Uri = "/api/admin/endpoints/keys/key-codex/codex-reset-credit/consume"
|
||||
.parse()
|
||||
.expect("uri should parse");
|
||||
let decision = classify_control_route(&http::Method::POST, &uri, &headers)
|
||||
.expect("decision should resolve");
|
||||
assert_eq!(decision.route_class.as_deref(), Some("admin_proxy"));
|
||||
assert_eq!(decision.route_family.as_deref(), Some("endpoints_manage"));
|
||||
assert_eq!(
|
||||
decision.route_kind.as_deref(),
|
||||
Some("codex_reset_credit_consume")
|
||||
);
|
||||
assert!(!decision.is_execution_runtime_candidate());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn admin_refresh_provider_quota_buffers_request_body_for_key_selection() {
|
||||
let headers = headers(&[]);
|
||||
|
||||
+116
@@ -0,0 +1,116 @@
|
||||
use crate::handlers::admin::provider::oauth::quota::codex::consume_codex_reset_credit_locally;
|
||||
use crate::handlers::admin::provider::oauth::quota::shared::{
|
||||
provider_quota_refresh_endpoint_for_provider, provider_quota_refresh_missing_endpoint_message,
|
||||
};
|
||||
use crate::handlers::admin::provider::shared::paths::admin_codex_reset_credit_consume_key_id;
|
||||
use crate::handlers::admin::provider::shared::payloads::AdminCodexResetCreditConsumeRequest;
|
||||
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
||||
use crate::GatewayError;
|
||||
use axum::{
|
||||
body::{Body, Bytes},
|
||||
http,
|
||||
response::{IntoResponse, Response},
|
||||
Json,
|
||||
};
|
||||
use serde_json::json;
|
||||
|
||||
pub(super) async fn maybe_handle(
|
||||
state: &AdminAppState<'_>,
|
||||
request_context: &AdminRequestContext<'_>,
|
||||
request_body: Option<&Bytes>,
|
||||
) -> Result<Option<Response<Body>>, GatewayError> {
|
||||
let Some(decision) = request_context.decision() else {
|
||||
return Ok(None);
|
||||
};
|
||||
if decision.route_family.as_deref() != Some("endpoints_manage")
|
||||
|| decision.route_kind.as_deref() != Some("codex_reset_credit_consume")
|
||||
|| request_context.method() != http::Method::POST
|
||||
|| !request_context
|
||||
.path()
|
||||
.starts_with("/api/admin/endpoints/keys/")
|
||||
|| !request_context
|
||||
.path()
|
||||
.ends_with("/codex-reset-credit/consume")
|
||||
{
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let Some(key_id) = admin_codex_reset_credit_consume_key_id(request_context.path()) else {
|
||||
return Ok(Some(not_found_response("Key 不存在")));
|
||||
};
|
||||
let payload = match request_body.filter(|body| !body.is_empty()) {
|
||||
Some(request_body) => {
|
||||
match serde_json::from_slice::<AdminCodexResetCreditConsumeRequest>(request_body) {
|
||||
Ok(payload) => payload,
|
||||
Err(_) => {
|
||||
return Ok(Some(bad_request_response("请求体必须是合法的 JSON 对象")));
|
||||
}
|
||||
}
|
||||
}
|
||||
None => {
|
||||
return Ok(Some(bad_request_response("请求体必须包含 idempotency_key")));
|
||||
}
|
||||
};
|
||||
let idempotency_key = payload.idempotency_key.trim().to_string();
|
||||
if idempotency_key.is_empty() {
|
||||
return Ok(Some(bad_request_response("idempotency_key 不能为空")));
|
||||
}
|
||||
|
||||
let Some(key) = state
|
||||
.read_provider_catalog_keys_by_ids(std::slice::from_ref(&key_id))
|
||||
.await?
|
||||
.into_iter()
|
||||
.next()
|
||||
else {
|
||||
return Ok(Some(not_found_response(format!("Key {key_id} 不存在"))));
|
||||
};
|
||||
let Some(provider) = state
|
||||
.read_provider_catalog_providers_by_ids(std::slice::from_ref(&key.provider_id))
|
||||
.await?
|
||||
.into_iter()
|
||||
.next()
|
||||
else {
|
||||
return Ok(Some(not_found_response(format!(
|
||||
"Provider {} 不存在",
|
||||
key.provider_id
|
||||
))));
|
||||
};
|
||||
let normalized_provider_type = provider.provider_type.trim().to_ascii_lowercase();
|
||||
if normalized_provider_type != "codex" {
|
||||
return Ok(Some(bad_request_response(
|
||||
"仅 Codex Provider 支持使用重置机会",
|
||||
)));
|
||||
}
|
||||
|
||||
let endpoints = state
|
||||
.list_provider_catalog_endpoints_by_provider_ids(std::slice::from_ref(&provider.id))
|
||||
.await?;
|
||||
let Some(endpoint) =
|
||||
provider_quota_refresh_endpoint_for_provider(&normalized_provider_type, &endpoints, true)
|
||||
else {
|
||||
return Ok(Some(bad_request_response(
|
||||
provider_quota_refresh_missing_endpoint_message(&normalized_provider_type),
|
||||
)));
|
||||
};
|
||||
|
||||
let (status, payload) =
|
||||
consume_codex_reset_credit_locally(state, &provider, &endpoint, key, &idempotency_key)
|
||||
.await?;
|
||||
Ok(Some((status, Json(payload)).into_response()))
|
||||
}
|
||||
|
||||
fn bad_request_response(detail: impl Into<String>) -> Response<Body> {
|
||||
(
|
||||
http::StatusCode::BAD_REQUEST,
|
||||
Json(json!({ "detail": detail.into() })),
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
|
||||
fn not_found_response(detail: impl Into<String>) -> Response<Body> {
|
||||
(
|
||||
http::StatusCode::NOT_FOUND,
|
||||
Json(json!({ "detail": detail.into() })),
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
mod batch;
|
||||
mod codex_reset_credit;
|
||||
mod create;
|
||||
mod delete;
|
||||
mod oauth_invalid;
|
||||
@@ -36,6 +37,11 @@ pub(super) async fn maybe_handle(
|
||||
{
|
||||
return Ok(Some(response));
|
||||
}
|
||||
if let Some(response) =
|
||||
codex_reset_credit::maybe_handle(state, request_context, request_body).await?
|
||||
{
|
||||
return Ok(Some(response));
|
||||
}
|
||||
if let Some(response) = create::maybe_handle(state, request_context, request_body).await? {
|
||||
return Ok(Some(response));
|
||||
}
|
||||
|
||||
@@ -8,10 +8,15 @@ use self::invalid::{
|
||||
codex_structured_invalid_reason,
|
||||
};
|
||||
use self::parse::{
|
||||
build_codex_quota_exhausted_fallback_metadata, parse_codex_usage_headers,
|
||||
build_codex_quota_exhausted_fallback_metadata, normalize_codex_reset_credit_consume_outcome,
|
||||
parse_codex_usage_headers, parse_codex_wham_reset_credits_detail_response,
|
||||
parse_codex_wham_usage_response,
|
||||
};
|
||||
use self::plan::{build_codex_quota_request_spec, execute_codex_quota_plan};
|
||||
use self::plan::{
|
||||
build_codex_quota_request_spec, build_codex_reset_credit_consume_request_spec,
|
||||
build_codex_reset_credits_request_spec, execute_codex_quota_plan,
|
||||
execute_codex_reset_credit_plan,
|
||||
};
|
||||
use super::shared::{
|
||||
build_quota_snapshot_payload, extract_execution_error_message,
|
||||
oauth_refresh_auto_removed_result, persist_provider_quota_refresh_state,
|
||||
@@ -26,7 +31,8 @@ use aether_contracts::ProxySnapshot;
|
||||
use aether_data_contracts::repository::provider_catalog::{
|
||||
StoredProviderCatalogEndpoint, StoredProviderCatalogKey, StoredProviderCatalogProvider,
|
||||
};
|
||||
use serde_json::json;
|
||||
use axum::http::StatusCode;
|
||||
use serde_json::{json, Map, Value};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
fn merge_codex_quota_metadata(
|
||||
@@ -45,6 +51,139 @@ fn merge_codex_quota_metadata(
|
||||
serde_json::Value::Object(merged)
|
||||
}
|
||||
|
||||
fn codex_reset_credits_available_count(metadata: &Map<String, Value>) -> Option<u64> {
|
||||
metadata
|
||||
.get("reset_credits")
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|reset_credits| reset_credits.get("available_count"))
|
||||
.and_then(aether_admin::provider::quota::coerce_json_u64)
|
||||
}
|
||||
|
||||
fn truncate_codex_reset_credit_detail_error(message: impl Into<String>) -> String {
|
||||
let message = message.into();
|
||||
let mut sanitized = message.replace('\n', " ");
|
||||
if sanitized.len() > 240 {
|
||||
sanitized.truncate(240);
|
||||
sanitized.push('…');
|
||||
}
|
||||
sanitized
|
||||
}
|
||||
|
||||
fn merge_codex_reset_credit_detail_metadata(
|
||||
codex_metadata: &mut Map<String, Value>,
|
||||
detail_metadata: &Value,
|
||||
) {
|
||||
let Some(detail_reset_credits) = detail_metadata
|
||||
.get("reset_credits")
|
||||
.and_then(Value::as_object)
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let mut reset_credits = codex_metadata
|
||||
.get("reset_credits")
|
||||
.and_then(Value::as_object)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
|
||||
let has_usage_available_count = reset_credits.contains_key("available_count");
|
||||
for (key, value) in detail_reset_credits {
|
||||
if key == "available_count" && has_usage_available_count {
|
||||
continue;
|
||||
}
|
||||
reset_credits.insert(key.clone(), value.clone());
|
||||
}
|
||||
codex_metadata.insert("reset_credits".to_string(), Value::Object(reset_credits));
|
||||
}
|
||||
|
||||
fn mark_codex_reset_credit_detail_failed(
|
||||
codex_metadata: &mut Map<String, Value>,
|
||||
detail_error: impl Into<String>,
|
||||
) {
|
||||
let mut reset_credits = codex_metadata
|
||||
.get("reset_credits")
|
||||
.and_then(Value::as_object)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
reset_credits.insert("detail_source".to_string(), json!("wham_readonly"));
|
||||
reset_credits.insert("detail_status".to_string(), json!("failed"));
|
||||
reset_credits.insert(
|
||||
"detail_error".to_string(),
|
||||
json!(truncate_codex_reset_credit_detail_error(detail_error)),
|
||||
);
|
||||
reset_credits
|
||||
.entry("credits".to_string())
|
||||
.or_insert_with(|| json!([]));
|
||||
codex_metadata.insert("reset_credits".to_string(), Value::Object(reset_credits));
|
||||
}
|
||||
|
||||
async fn enrich_codex_reset_credit_details(
|
||||
state: &AdminAppState<'_>,
|
||||
transport: &crate::handlers::admin::request::AdminGatewayProviderTransportSnapshot,
|
||||
resolved_oauth_auth: Option<(String, String)>,
|
||||
proxy_override: Option<&ProxySnapshot>,
|
||||
codex_metadata: &mut Map<String, Value>,
|
||||
now_unix_secs: u64,
|
||||
) -> Result<(), GatewayError> {
|
||||
let available_count = codex_reset_credits_available_count(codex_metadata).unwrap_or(0);
|
||||
if available_count == 0 {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let request_spec = match build_codex_reset_credits_request_spec(transport, resolved_oauth_auth)
|
||||
{
|
||||
Ok(request_spec) => request_spec,
|
||||
Err(message) => {
|
||||
mark_codex_reset_credit_detail_failed(codex_metadata, message);
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
|
||||
let result =
|
||||
match execute_codex_reset_credit_plan(state, transport, request_spec, proxy_override)
|
||||
.await?
|
||||
{
|
||||
ProviderQuotaExecutionOutcome::Response(result) => result,
|
||||
ProviderQuotaExecutionOutcome::Failure(detail) => {
|
||||
mark_codex_reset_credit_detail_failed(
|
||||
codex_metadata,
|
||||
format!("reset credit detail 请求执行失败: {detail}"),
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
|
||||
if result.status_code != 200 {
|
||||
let detail = extract_execution_error_message(&result)
|
||||
.unwrap_or_else(|| format!("HTTP {}", result.status_code));
|
||||
mark_codex_reset_credit_detail_failed(
|
||||
codex_metadata,
|
||||
format!(
|
||||
"reset credit detail 返回状态码 {}: {detail}",
|
||||
result.status_code
|
||||
),
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let Some(body_json) = result
|
||||
.body
|
||||
.as_ref()
|
||||
.and_then(|body| body.json_body.as_ref())
|
||||
else {
|
||||
mark_codex_reset_credit_detail_failed(codex_metadata, "无法解析 reset credit detail 响应");
|
||||
return Ok(());
|
||||
};
|
||||
if let Some(detail_metadata) =
|
||||
parse_codex_wham_reset_credits_detail_response(body_json, now_unix_secs)
|
||||
{
|
||||
merge_codex_reset_credit_detail_metadata(codex_metadata, &detail_metadata);
|
||||
} else {
|
||||
mark_codex_reset_credit_detail_failed(codex_metadata, "reset credit detail 响应为空");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn codex_oauth_refresh_issue_reason(reason: Option<&str>) -> bool {
|
||||
reason.is_some_and(|reason| {
|
||||
reason
|
||||
@@ -54,6 +193,210 @@ fn codex_oauth_refresh_issue_reason(reason: Option<&str>) -> bool {
|
||||
})
|
||||
}
|
||||
|
||||
fn codex_consume_success_status(outcome: &str) -> &'static str {
|
||||
match outcome {
|
||||
"reset" | "already_redeemed" => "success",
|
||||
"nothing_to_reset" | "no_credit" => "noop",
|
||||
_ => "unknown",
|
||||
}
|
||||
}
|
||||
|
||||
fn codex_extract_refresh_result_fields(
|
||||
refresh_payload: Option<&Value>,
|
||||
key_id: &str,
|
||||
) -> (String, Option<String>, Option<Value>, Option<Value>) {
|
||||
let Some(result) = refresh_payload
|
||||
.and_then(|payload| payload.get("results"))
|
||||
.and_then(Value::as_array)
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.filter_map(Value::as_object)
|
||||
.find(|item| item.get("key_id").and_then(Value::as_str) == Some(key_id))
|
||||
else {
|
||||
return (
|
||||
"failed".to_string(),
|
||||
Some("刷新结果中缺少当前 key".to_string()),
|
||||
None,
|
||||
None,
|
||||
);
|
||||
};
|
||||
|
||||
let status = result
|
||||
.get("status")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.unwrap_or_default();
|
||||
let refresh_status = if status.eq_ignore_ascii_case("success") {
|
||||
"success"
|
||||
} else {
|
||||
"failed"
|
||||
}
|
||||
.to_string();
|
||||
let refresh_error = if refresh_status == "success" {
|
||||
None
|
||||
} else {
|
||||
result
|
||||
.get("message")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.map(ToOwned::to_owned)
|
||||
};
|
||||
(
|
||||
refresh_status,
|
||||
refresh_error,
|
||||
result.get("metadata").cloned(),
|
||||
result.get("quota_snapshot").cloned(),
|
||||
)
|
||||
}
|
||||
|
||||
pub(crate) async fn consume_codex_reset_credit_locally(
|
||||
state: &AdminAppState<'_>,
|
||||
provider: &StoredProviderCatalogProvider,
|
||||
endpoint: &StoredProviderCatalogEndpoint,
|
||||
key: StoredProviderCatalogKey,
|
||||
idempotency_key: &str,
|
||||
) -> Result<(StatusCode, Value), GatewayError> {
|
||||
let transport = match state
|
||||
.read_provider_transport_snapshot(&provider.id, &endpoint.id, &key.id)
|
||||
.await?
|
||||
{
|
||||
Some(transport) => transport,
|
||||
None => {
|
||||
return Ok((
|
||||
StatusCode::BAD_GATEWAY,
|
||||
json!({
|
||||
"key_id": key.id,
|
||||
"status": "error",
|
||||
"outcome": "error",
|
||||
"message": "Provider transport snapshot unavailable",
|
||||
}),
|
||||
));
|
||||
}
|
||||
};
|
||||
|
||||
let is_oauth_managed = provider_key_is_oauth_managed(&key, provider.provider_type.as_str());
|
||||
let resolved_oauth_auth = if is_oauth_managed {
|
||||
state.resolve_local_oauth_header_auth(&transport).await?
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if is_oauth_managed && resolved_oauth_auth.is_none() {
|
||||
return Ok((
|
||||
StatusCode::BAD_REQUEST,
|
||||
json!({
|
||||
"key_id": key.id,
|
||||
"status": "error",
|
||||
"outcome": "error",
|
||||
"message": "缺少 Codex OAuth 认证信息,请先重新授权/刷新 Token",
|
||||
}),
|
||||
));
|
||||
}
|
||||
|
||||
let request_spec = match build_codex_reset_credit_consume_request_spec(
|
||||
&transport,
|
||||
resolved_oauth_auth,
|
||||
idempotency_key,
|
||||
) {
|
||||
Ok(request_spec) => request_spec,
|
||||
Err(message) => {
|
||||
return Ok((
|
||||
StatusCode::BAD_REQUEST,
|
||||
json!({
|
||||
"key_id": key.id,
|
||||
"status": "error",
|
||||
"outcome": "error",
|
||||
"message": message,
|
||||
}),
|
||||
));
|
||||
}
|
||||
};
|
||||
|
||||
let result =
|
||||
match execute_codex_reset_credit_plan(state, &transport, request_spec, None).await? {
|
||||
ProviderQuotaExecutionOutcome::Response(result) => result,
|
||||
ProviderQuotaExecutionOutcome::Failure(detail) => {
|
||||
return Ok((
|
||||
StatusCode::BAD_GATEWAY,
|
||||
json!({
|
||||
"key_id": key.id,
|
||||
"status": "error",
|
||||
"outcome": "error",
|
||||
"message": format!("reset credit consume 请求执行失败: {detail}"),
|
||||
}),
|
||||
));
|
||||
}
|
||||
};
|
||||
|
||||
let body_json = result
|
||||
.body
|
||||
.as_ref()
|
||||
.and_then(|body| body.json_body.as_ref());
|
||||
let outcome = normalize_codex_reset_credit_consume_outcome(body_json)
|
||||
.unwrap_or_else(|| "unknown".to_string());
|
||||
let known_non_error_outcome = matches!(
|
||||
outcome.as_str(),
|
||||
"reset" | "already_redeemed" | "nothing_to_reset" | "no_credit"
|
||||
);
|
||||
if result.status_code >= 400 && !known_non_error_outcome {
|
||||
let detail = extract_execution_error_message(&result)
|
||||
.unwrap_or_else(|| format!("HTTP {}", result.status_code));
|
||||
return Ok((
|
||||
StatusCode::BAD_GATEWAY,
|
||||
json!({
|
||||
"key_id": key.id,
|
||||
"status": "error",
|
||||
"outcome": "error",
|
||||
"idempotency_key": idempotency_key,
|
||||
"message": format!("reset credit consume 返回状态码 {}: {detail}", result.status_code),
|
||||
"status_code": result.status_code,
|
||||
}),
|
||||
));
|
||||
}
|
||||
|
||||
let (refresh_status, refresh_error, metadata, quota_snapshot) =
|
||||
match refresh_codex_provider_quota_locally(
|
||||
state,
|
||||
provider,
|
||||
endpoint,
|
||||
vec![key.clone()],
|
||||
None,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(refresh_payload) => {
|
||||
codex_extract_refresh_result_fields(refresh_payload.as_ref(), &key.id)
|
||||
}
|
||||
Err(err) => (
|
||||
"failed".to_string(),
|
||||
Some(truncate_codex_reset_credit_detail_error(err.into_message())),
|
||||
None,
|
||||
None,
|
||||
),
|
||||
};
|
||||
|
||||
let mut payload = Map::new();
|
||||
payload.insert("key_id".to_string(), json!(key.id));
|
||||
payload.insert(
|
||||
"status".to_string(),
|
||||
json!(codex_consume_success_status(&outcome)),
|
||||
);
|
||||
payload.insert("outcome".to_string(), json!(outcome));
|
||||
payload.insert("idempotency_key".to_string(), json!(idempotency_key));
|
||||
payload.insert("refresh_status".to_string(), json!(refresh_status));
|
||||
if let Some(refresh_error) = refresh_error {
|
||||
payload.insert("refresh_error".to_string(), json!(refresh_error));
|
||||
}
|
||||
if let Some(metadata) = metadata {
|
||||
payload.insert("metadata".to_string(), metadata);
|
||||
}
|
||||
if let Some(quota_snapshot) = quota_snapshot {
|
||||
payload.insert("quota_snapshot".to_string(), quota_snapshot);
|
||||
}
|
||||
|
||||
Ok((StatusCode::OK, Value::Object(payload)))
|
||||
}
|
||||
|
||||
pub(crate) async fn refresh_codex_provider_quota_locally(
|
||||
state: &AdminAppState<'_>,
|
||||
provider: &StoredProviderCatalogProvider,
|
||||
@@ -114,19 +457,20 @@ pub(crate) async fn refresh_codex_provider_quota_locally(
|
||||
continue;
|
||||
}
|
||||
|
||||
let request_spec = match build_codex_quota_request_spec(&transport, resolved_oauth_auth) {
|
||||
Ok(request_spec) => request_spec,
|
||||
Err(message) => {
|
||||
failed_count += 1;
|
||||
results.push(json!({
|
||||
"key_id": key.id,
|
||||
"key_name": key.name,
|
||||
"status": "error",
|
||||
"message": message,
|
||||
}));
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let request_spec =
|
||||
match build_codex_quota_request_spec(&transport, resolved_oauth_auth.clone()) {
|
||||
Ok(request_spec) => request_spec,
|
||||
Err(message) => {
|
||||
failed_count += 1;
|
||||
results.push(json!({
|
||||
"key_id": key.id,
|
||||
"key_name": key.name,
|
||||
"status": "error",
|
||||
"message": message,
|
||||
}));
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let result = match execute_codex_quota_plan(
|
||||
state,
|
||||
@@ -171,8 +515,22 @@ pub(crate) async fn refresh_codex_provider_quota_locally(
|
||||
.and_then(|body| body.json_body.as_ref())
|
||||
{
|
||||
if let Some(parsed) = parse_codex_wham_usage_response(body_json, now_unix_secs) {
|
||||
let mut codex_metadata =
|
||||
match merge_codex_quota_metadata(header_metadata.as_ref(), &parsed) {
|
||||
Value::Object(object) => object,
|
||||
_ => Map::new(),
|
||||
};
|
||||
enrich_codex_reset_credit_details(
|
||||
state,
|
||||
&transport,
|
||||
resolved_oauth_auth.clone(),
|
||||
proxy_override.as_ref(),
|
||||
&mut codex_metadata,
|
||||
now_unix_secs,
|
||||
)
|
||||
.await?;
|
||||
metadata_update = Some(json!({
|
||||
"codex": merge_codex_quota_metadata(header_metadata.as_ref(), &parsed)
|
||||
"codex": codex_metadata
|
||||
}));
|
||||
(oauth_invalid_at_unix_secs, oauth_invalid_reason) =
|
||||
quota_refresh_success_invalid_state(&key);
|
||||
|
||||
@@ -22,6 +22,22 @@ pub(super) fn parse_codex_wham_usage_response(
|
||||
admin_provider_quota_pure::parse_codex_wham_usage_response(value, updated_at_unix_secs)
|
||||
}
|
||||
|
||||
pub(super) fn parse_codex_wham_reset_credits_detail_response(
|
||||
value: &serde_json::Value,
|
||||
updated_at_unix_secs: u64,
|
||||
) -> Option<serde_json::Value> {
|
||||
admin_provider_quota_pure::parse_codex_wham_reset_credits_detail_response(
|
||||
value,
|
||||
updated_at_unix_secs,
|
||||
)
|
||||
}
|
||||
|
||||
pub(super) fn normalize_codex_reset_credit_consume_outcome(
|
||||
value: Option<&serde_json::Value>,
|
||||
) -> Option<String> {
|
||||
admin_provider_quota_pure::normalize_codex_reset_credit_consume_outcome(value)
|
||||
}
|
||||
|
||||
pub(super) fn parse_codex_usage_headers(
|
||||
headers: &BTreeMap<String, String>,
|
||||
updated_at_unix_secs: u64,
|
||||
|
||||
@@ -5,17 +5,26 @@ use super::super::shared::{
|
||||
use crate::handlers::admin::request::{AdminAppState, AdminGatewayProviderTransportSnapshot};
|
||||
use crate::GatewayError;
|
||||
use aether_contracts::ProxySnapshot;
|
||||
use aether_provider_pool::{build_codex_pool_quota_request, ProviderPoolQuotaRequestSpec};
|
||||
use aether_provider_pool::{
|
||||
build_codex_pool_quota_request, build_codex_pool_reset_credit_consume_request,
|
||||
build_codex_pool_reset_credits_request, ProviderPoolQuotaRequestSpec,
|
||||
};
|
||||
|
||||
fn codex_auth_config(
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
) -> Option<serde_json::Value> {
|
||||
transport
|
||||
.key
|
||||
.decrypted_auth_config
|
||||
.as_deref()
|
||||
.and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
|
||||
}
|
||||
|
||||
pub(super) fn build_codex_quota_request_spec(
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
resolved_oauth_auth: Option<(String, String)>,
|
||||
) -> Result<ProviderPoolQuotaRequestSpec, String> {
|
||||
let auth_config = transport
|
||||
.key
|
||||
.decrypted_auth_config
|
||||
.as_deref()
|
||||
.and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok());
|
||||
let auth_config = codex_auth_config(transport);
|
||||
let mut request = build_codex_pool_quota_request(
|
||||
&transport.key.id,
|
||||
resolved_oauth_auth,
|
||||
@@ -29,6 +38,44 @@ pub(super) fn build_codex_quota_request_spec(
|
||||
Ok(request)
|
||||
}
|
||||
|
||||
pub(super) fn build_codex_reset_credits_request_spec(
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
resolved_oauth_auth: Option<(String, String)>,
|
||||
) -> Result<ProviderPoolQuotaRequestSpec, String> {
|
||||
let auth_config = codex_auth_config(transport);
|
||||
let mut request = build_codex_pool_reset_credits_request(
|
||||
&transport.key.id,
|
||||
resolved_oauth_auth,
|
||||
Some(transport.key.decrypted_api_key.as_str()),
|
||||
auth_config.as_ref(),
|
||||
)?;
|
||||
crate::provider_transport::apply_local_auth_config_header_overrides(
|
||||
&mut request.headers,
|
||||
transport.key.decrypted_auth_config.as_deref(),
|
||||
);
|
||||
Ok(request)
|
||||
}
|
||||
|
||||
pub(super) fn build_codex_reset_credit_consume_request_spec(
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
resolved_oauth_auth: Option<(String, String)>,
|
||||
idempotency_key: &str,
|
||||
) -> Result<ProviderPoolQuotaRequestSpec, String> {
|
||||
let auth_config = codex_auth_config(transport);
|
||||
let mut request = build_codex_pool_reset_credit_consume_request(
|
||||
&transport.key.id,
|
||||
resolved_oauth_auth,
|
||||
Some(transport.key.decrypted_api_key.as_str()),
|
||||
auth_config.as_ref(),
|
||||
idempotency_key,
|
||||
)?;
|
||||
crate::provider_transport::apply_local_auth_config_header_overrides(
|
||||
&mut request.headers,
|
||||
transport.key.decrypted_auth_config.as_deref(),
|
||||
);
|
||||
Ok(request)
|
||||
}
|
||||
|
||||
pub(super) async fn execute_codex_quota_plan(
|
||||
state: &AdminAppState<'_>,
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
@@ -56,3 +103,31 @@ pub(super) async fn execute_codex_quota_plan(
|
||||
);
|
||||
execute_provider_quota_plan(state, transport, plan, "codex").await
|
||||
}
|
||||
|
||||
pub(super) async fn execute_codex_reset_credit_plan(
|
||||
state: &AdminAppState<'_>,
|
||||
transport: &AdminGatewayProviderTransportSnapshot,
|
||||
spec: ProviderPoolQuotaRequestSpec,
|
||||
proxy_override: Option<&ProxySnapshot>,
|
||||
) -> Result<ProviderQuotaExecutionOutcome, GatewayError> {
|
||||
let proxy = match proxy_override {
|
||||
Some(proxy) => Some(proxy.clone()),
|
||||
None => {
|
||||
state
|
||||
.resolve_transport_proxy_snapshot_with_tunnel_affinity(transport)
|
||||
.await
|
||||
}
|
||||
};
|
||||
let timeouts = Some(resolve_provider_quota_execution_timeouts(
|
||||
state.resolve_transport_execution_timeouts(transport),
|
||||
proxy.as_ref(),
|
||||
));
|
||||
let plan = build_provider_quota_execution_plan(
|
||||
transport,
|
||||
spec,
|
||||
proxy,
|
||||
state.resolve_transport_profile(transport),
|
||||
timeouts,
|
||||
);
|
||||
execute_provider_quota_plan(state, transport, plan, "codex_reset_credit").await
|
||||
}
|
||||
|
||||
@@ -40,6 +40,13 @@ pub(crate) fn admin_reset_cycle_stats_key_id(request_path: &str) -> Option<Strin
|
||||
.map(ToOwned::to_owned)
|
||||
}
|
||||
|
||||
pub(crate) fn admin_codex_reset_credit_consume_key_id(request_path: &str) -> Option<String> {
|
||||
request_path
|
||||
.strip_prefix("/api/admin/endpoints/keys/")?
|
||||
.strip_suffix("/codex-reset-credit/consume")
|
||||
.map(ToOwned::to_owned)
|
||||
}
|
||||
|
||||
pub(crate) fn admin_update_key_id(request_path: &str) -> Option<String> {
|
||||
let key_id = request_path.strip_prefix("/api/admin/endpoints/keys/")?;
|
||||
(!key_id.is_empty() && !key_id.contains('/')).then_some(key_id.to_string())
|
||||
|
||||
@@ -15,9 +15,9 @@ pub(crate) use self::crud::{
|
||||
is_admin_providers_root,
|
||||
};
|
||||
pub(crate) use self::endpoint_keys::{
|
||||
admin_clear_oauth_invalid_key_id, admin_export_key_id, admin_provider_id_for_keys,
|
||||
admin_provider_id_for_refresh_quota, admin_reset_cycle_stats_key_id, admin_reveal_key_id,
|
||||
admin_update_key_id,
|
||||
admin_clear_oauth_invalid_key_id, admin_codex_reset_credit_consume_key_id, admin_export_key_id,
|
||||
admin_provider_id_for_keys, admin_provider_id_for_refresh_quota,
|
||||
admin_reset_cycle_stats_key_id, admin_reveal_key_id, admin_update_key_id,
|
||||
};
|
||||
pub(crate) use self::oauth::{
|
||||
admin_provider_oauth_batch_import_provider_id, admin_provider_oauth_batch_import_task_path,
|
||||
|
||||
@@ -113,6 +113,11 @@ pub(crate) struct AdminProviderQuotaRefreshRequest {
|
||||
pub(crate) key_ids: Option<Vec<String>>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub(crate) struct AdminCodexResetCreditConsumeRequest {
|
||||
pub(crate) idempotency_key: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub(crate) struct AdminProviderCreateRequest {
|
||||
pub(crate) name: String,
|
||||
|
||||
@@ -855,6 +855,7 @@ fn build_codex_quota_status_snapshot(
|
||||
let credits_unlimited = metadata
|
||||
.get("credits_unlimited")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_bool);
|
||||
let reset_credits = build_codex_reset_credits_status_snapshot(metadata, observed_at_unix_secs);
|
||||
|
||||
let windows = [
|
||||
codex_quota_window_snapshot(metadata, "primary", "weekly", "周", observed_at_unix_secs),
|
||||
@@ -883,6 +884,7 @@ fn build_codex_quota_status_snapshot(
|
||||
&& credits_has_credits.is_none()
|
||||
&& credits_balance.is_none()
|
||||
&& credits_unlimited.is_none()
|
||||
&& reset_credits.is_none()
|
||||
&& observed_at_unix_secs.is_none()
|
||||
{
|
||||
return None;
|
||||
@@ -958,6 +960,7 @@ fn build_codex_quota_status_snapshot(
|
||||
} else {
|
||||
Value::Object(credits)
|
||||
},
|
||||
"reset_credits": reset_credits,
|
||||
"windows": windows,
|
||||
}))
|
||||
}
|
||||
@@ -1815,6 +1818,119 @@ fn build_gemini_cli_quota_status_snapshot(
|
||||
}))
|
||||
}
|
||||
|
||||
fn build_codex_reset_credits_status_snapshot(
|
||||
metadata: &Map<String, Value>,
|
||||
observed_at_unix_secs: Option<u64>,
|
||||
) -> Option<Value> {
|
||||
let reset_credits = metadata.get("reset_credits").and_then(Value::as_object)?;
|
||||
let available_count = reset_credits
|
||||
.get("available_count")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_u64);
|
||||
let updated_at = reset_credits
|
||||
.get("updated_at")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_u64)
|
||||
.or(observed_at_unix_secs);
|
||||
let detail_source = reset_credits
|
||||
.get("detail_source")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty());
|
||||
let detail_status = reset_credits
|
||||
.get("detail_status")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty());
|
||||
let detail_error = reset_credits
|
||||
.get("detail_error")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty());
|
||||
|
||||
let mut credits = reset_credits
|
||||
.get("credits")
|
||||
.and_then(Value::as_array)
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.filter_map(|item| {
|
||||
let object = item.as_object()?;
|
||||
let expires_at = object
|
||||
.get("expires_at")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_u64)?;
|
||||
let display_key = object
|
||||
.get("display_key")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())?;
|
||||
let mut out = Map::new();
|
||||
if let Some(id) = object
|
||||
.get("id")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
out.insert("id".to_string(), json!(id));
|
||||
}
|
||||
out.insert("display_key".to_string(), json!(display_key));
|
||||
if let Some(status) = object
|
||||
.get("status")
|
||||
.and_then(Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
out.insert("status".to_string(), json!(status));
|
||||
}
|
||||
if let Some(granted_at) = object
|
||||
.get("granted_at")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_u64)
|
||||
{
|
||||
out.insert("granted_at".to_string(), json!(granted_at));
|
||||
}
|
||||
out.insert("expires_at".to_string(), json!(expires_at));
|
||||
if let Some(observed_at) = observed_at_unix_secs {
|
||||
out.insert(
|
||||
"remaining_seconds".to_string(),
|
||||
json!(expires_at.saturating_sub(observed_at)),
|
||||
);
|
||||
}
|
||||
Some(Value::Object(out))
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
credits.sort_by_key(|item| {
|
||||
item.get("expires_at")
|
||||
.and_then(admin_provider_quota_pure::coerce_json_u64)
|
||||
.unwrap_or(u64::MAX)
|
||||
});
|
||||
|
||||
if available_count.is_none()
|
||||
&& updated_at.is_none()
|
||||
&& detail_source.is_none()
|
||||
&& detail_status.is_none()
|
||||
&& detail_error.is_none()
|
||||
&& credits.is_empty()
|
||||
{
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut out = Map::new();
|
||||
if let Some(value) = available_count {
|
||||
out.insert("available_count".to_string(), json!(value));
|
||||
}
|
||||
if let Some(value) = updated_at {
|
||||
out.insert("updated_at".to_string(), json!(value));
|
||||
}
|
||||
if let Some(value) = detail_source {
|
||||
out.insert("detail_source".to_string(), json!(value));
|
||||
}
|
||||
if let Some(value) = detail_status {
|
||||
out.insert("detail_status".to_string(), json!(value));
|
||||
}
|
||||
if let Some(value) = detail_error {
|
||||
out.insert("detail_error".to_string(), json!(value));
|
||||
}
|
||||
out.insert("credits".to_string(), Value::Array(credits));
|
||||
Some(Value::Object(out))
|
||||
}
|
||||
|
||||
pub(crate) fn sync_provider_key_quota_status_snapshot(
|
||||
status_snapshot: Option<&Value>,
|
||||
provider_type: &str,
|
||||
@@ -1882,6 +1998,12 @@ fn quota_snapshot_has_materialized_data(
|
||||
{
|
||||
return true;
|
||||
}
|
||||
if quota_snapshot
|
||||
.get("reset_credits")
|
||||
.is_some_and(|reset_credits| !reset_credits.is_null())
|
||||
{
|
||||
return true;
|
||||
}
|
||||
|
||||
quota_snapshot
|
||||
.get("code")
|
||||
@@ -2644,6 +2766,68 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn provider_key_status_snapshot_payload_backfills_codex_reset_credits() {
|
||||
let mut key = sample_catalog_key();
|
||||
key.upstream_metadata = Some(json!({
|
||||
"codex": {
|
||||
"updated_at": 1_775_553_285u64,
|
||||
"plan_type": "plus",
|
||||
"primary_used_percent": 55.0,
|
||||
"primary_reset_at": 1_900_000_000u64,
|
||||
"has_credits": true,
|
||||
"credits_balance": 42.0,
|
||||
"reset_credits": {
|
||||
"available_count": 2,
|
||||
"updated_at": 1_775_553_285u64,
|
||||
"detail_source": "wham_readonly",
|
||||
"detail_status": "available",
|
||||
"credits": [
|
||||
{
|
||||
"id": "bbbbbbbb-1111-2222-3333-444444444444",
|
||||
"display_key": "bbbbbbbb",
|
||||
"status": "available",
|
||||
"expires_at": 1_775_900_000u64
|
||||
},
|
||||
{
|
||||
"id": "aaaaaaaa-1111-2222-3333-444444444444",
|
||||
"display_key": "aaaaaaaa",
|
||||
"status": "available",
|
||||
"expires_at": 1_775_700_000u64
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
}));
|
||||
|
||||
let payload = provider_key_status_snapshot_payload(&key, "codex");
|
||||
let quota = payload
|
||||
.get("quota")
|
||||
.and_then(Value::as_object)
|
||||
.expect("quota snapshot should be object");
|
||||
|
||||
assert_eq!(quota.get("exhausted"), Some(&json!(false)));
|
||||
assert_eq!(
|
||||
payload.pointer("/quota/reset_credits/available_count"),
|
||||
Some(&json!(2u64))
|
||||
);
|
||||
assert_eq!(
|
||||
payload.pointer("/quota/reset_credits/credits/0/display_key"),
|
||||
Some(&json!("aaaaaaaa"))
|
||||
);
|
||||
assert_eq!(
|
||||
payload.pointer("/quota/reset_credits/credits/0/remaining_seconds"),
|
||||
Some(&json!(146_715u64))
|
||||
);
|
||||
assert_eq!(
|
||||
quota
|
||||
.get("credits")
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|credits| credits.get("balance")),
|
||||
Some(&json!(42.0))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn provider_key_status_snapshot_payload_backfills_codex_spark_windows() {
|
||||
let mut key = sample_catalog_key();
|
||||
|
||||
Reference in New Issue
Block a user