Merge origin/main into main

This commit is contained in:
elky
2026-10-05 12:10:23 +08:00
22 changed files with 2050 additions and 50 deletions
@@ -1930,6 +1930,29 @@ mod tests {
assert_eq!(active["reasoning_effort"], "max");
}
#[test]
fn user_usage_payloads_expose_gemini_thinking_config_reasoning_mapping() {
let item = StoredRequestUsageAudit {
request_body: Some(json!({
"generationConfig": {
"thinkingConfig": { "includeThoughts": true, "thinkingLevel": "HIGH" }
}
})),
provider_request_body: Some(json!({
"generationConfig": { "thinkingConfig": { "thinkingBudget": 8192 } }
})),
..sample_usage("completed")
};
let record = build_users_me_usage_record_payload(&item, false, &BTreeMap::new(), false);
let active = build_users_me_usage_active_payload(&item);
assert_eq!(record["requested_reasoning_effort"], "high");
assert_eq!(active["requested_reasoning_effort"], "high");
assert_eq!(record["reasoning_effort"], "xhigh");
assert_eq!(active["reasoning_effort"], "xhigh");
}
#[test]
fn user_usage_payloads_expose_websocket_transport() {
let item = StoredRequestUsageAudit {
@@ -2681,6 +2681,84 @@ fn windsurf_quota_snapshot_has_stale_cooldown(quota_snapshot: &Map<String, Value
has_capacity || !exhausted
}
/// 读取状态快照时,把已越过重置时间点的配额窗口按“已重置”口径归一化。
///
/// 背景:调度侧早已把到期窗口视为未耗尽(`provider_pool_reset_deadline_elapsed`),
/// 账号额度文本也会按到期强制显示 100%,但列表读取此前直接返回存量快照,
/// 导致管理端倒计时归零后进度条仍停留在旧的剩余百分比。这里让读取层与
/// 调度侧、文本侧使用同一口径,避免三处状态互相矛盾。
///
/// 覆盖所有带重置时间的提供商窗口(codex/kiro/xai/grok/antigravity/
/// gemini_cli/chatgpt_web/windsurf 等):比例(used_ratio/remaining_ratio)、
/// 数值(used_value/remaining_value)与窗口级耗尽标记会一起恢复为“已重置”。
fn normalize_expired_quota_windows(snapshot: &mut serde_json::Map<String, Value>) {
let Some(quota) = snapshot.get_mut("quota").and_then(Value::as_object_mut) else {
return;
};
let now_unix_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()
.map(|duration| duration.as_secs())
.unwrap_or(0);
let fallback_observed_at = provider_quota_timestamp_unix_secs(quota.get("observed_at"))
.or_else(|| provider_quota_timestamp_unix_secs(quota.get("updated_at")));
let Some(windows) = quota.get_mut("windows").and_then(Value::as_array_mut) else {
return;
};
for window in windows.iter_mut().filter_map(Value::as_object_mut) {
// 与同类窗口处理保持一致:window_minutes=0 不是真实配额窗口。
if window.get("window_minutes").and_then(Value::as_u64) == Some(0) {
continue;
}
if !aether_provider_pool::provider_pool_reset_deadline_elapsed(
window,
fallback_observed_at,
now_unix_secs,
) {
continue;
}
// 只归一化“观测完整”的窗口:要么有比例观测,要么有上限 + 用量数值。
// 这样既不会把无数据窗口凭空显示成 100%,也不会出现比例已恢复 100%
// 而数值仍停留在旧值的不一致(例如只有 remaining_value 却没有上限的窗口)。
let has_ratio_observation = ["used_ratio", "remaining_ratio"]
.into_iter()
.any(|field| window.get(field).is_some_and(Value::is_number));
let limit_value = window
.get("limit_value")
.and_then(admin_provider_quota_pure::coerce_json_f64)
.filter(|value| *value > 0.0);
let has_value_observation = limit_value.is_some()
&& ["used_value", "remaining_value"]
.into_iter()
.any(|field| window.get(field).is_some_and(Value::is_number));
if !has_ratio_observation && !has_value_observation {
continue;
}
// 比例口径:已用清零、剩余 100%;字段原本为 null 时一并补齐,保证展示口径统一。
for (field, value) in [("used_ratio", 0.0), ("remaining_ratio", 1.0)] {
if window.contains_key(field) {
window.insert(field.to_string(), json!(value));
}
}
// 数值口径:已用清零;有上限时把剩余恢复到上限(用于“剩余 x/y”类文本展示)。
if window.contains_key("used_value") {
window.insert("used_value".to_string(), json!(0.0));
}
if window.contains_key("remaining_value") {
if let Some(limit_value) = limit_value {
window.insert("remaining_value".to_string(), json!(limit_value));
}
}
// 同步清掉窗口级耗尽标记,避免展示与调度口径互相矛盾。
for field in ["is_exhausted", "exhausted"] {
if let Some(slot) = window.get_mut(field) {
*slot = json!(false);
}
}
}
}
pub(crate) fn provider_key_status_snapshot_payload(
key: &StoredProviderCatalogKey,
provider_type: &str,
@@ -2722,9 +2800,8 @@ pub(crate) fn provider_key_status_snapshot_payload(
// Legacy snapshots can retain an exhausted summary after a window reset or
// newer quota observation. Use the same decision as scheduling so the
// account list and its status filter do not keep displaying that stale block.
if provider_type.trim().eq_ignore_ascii_case("codex")
&& !aether_provider_pool::provider_pool_key_account_quota_exhausted(key, provider_type)
{
// 对所有提供商生效:适配器判定已是“重置感知”的,与调度口径保持一致。
if !aether_provider_pool::provider_pool_key_account_quota_exhausted(key, provider_type) {
if let Some(quota) = snapshot.get_mut("quota").and_then(Value::as_object_mut) {
quota.insert("exhausted".to_string(), json!(false));
if quota.get("code").and_then(Value::as_str) == Some("exhausted") {
@@ -2732,6 +2809,9 @@ pub(crate) fn provider_key_status_snapshot_payload(
}
}
}
// 读取时归一化已到期的配额窗口,保证列表进度条、额度文字与调度侧、
// 账号额度文本使用同一“已重置”口径。
normalize_expired_quota_windows(&mut snapshot);
snapshot.insert(
"oauth".to_string(),
build_provider_key_oauth_status_snapshot(key),
@@ -3780,6 +3860,421 @@ mod tests {
assert_eq!(window.get("reset_seconds"), Some(&json!(3_600u64)));
}
#[test]
fn provider_key_status_snapshot_payload_normalizes_expired_codex_quota_windows() {
let mut key = sample_catalog_key();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "codex",
"code": "ok",
"exhausted": false,
"updated_at": 1_700_000_000u64,
"windows": [
{
"code": "5h",
"label": "5H",
"scope": "account",
"unit": "percent",
"used_ratio": 0.88,
"remaining_ratio": 0.12,
"reset_at": 1_700_003_600u64,
"window_minutes": 300
},
{
"code": "weekly",
"label": "周",
"scope": "account",
"unit": "percent",
"used_ratio": 0.5,
"remaining_ratio": 0.5,
"reset_at": 2_000_000_000u64,
"window_minutes": 10_080
},
{
"code": "spark_5h",
"label": "Spark 5H",
"scope": "account",
"unit": "percent",
"reset_at": 1_700_003_600u64,
"window_minutes": 300
},
{
"code": "unlimited",
"label": "无限",
"scope": "account",
"unit": "percent",
"used_ratio": 0.3,
"remaining_ratio": 0.7,
"reset_at": 1_700_003_600u64,
"window_minutes": 0
}
]
}
}));
let payload = provider_key_status_snapshot_payload(&key, "codex");
// 已到期的窗口按“已重置”口径归一化:用量清零、剩余 100%。
assert_eq!(payload.pointer("/quota/windows/0/code"), Some(&json!("5h")));
assert_eq!(
payload.pointer("/quota/windows/0/used_ratio"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_ratio"),
Some(&json!(1.0))
);
// 未到期的窗口保留原始观测值。
assert_eq!(
payload.pointer("/quota/windows/1/code"),
Some(&json!("weekly"))
);
assert_eq!(
payload.pointer("/quota/windows/1/used_ratio"),
Some(&json!(0.5))
);
assert_eq!(
payload.pointer("/quota/windows/1/remaining_ratio"),
Some(&json!(0.5))
);
// 没有用量观测的窗口不会被凭空补成 100%。
assert_eq!(
payload.pointer("/quota/windows/2/code"),
Some(&json!("spark_5h"))
);
assert_eq!(payload.pointer("/quota/windows/2/used_ratio"), None);
assert_eq!(payload.pointer("/quota/windows/2/remaining_ratio"), None);
// window_minutes=0 不是真实配额窗口,与同类 Codex 窗口处理保持一致,不归一化。
assert_eq!(
payload.pointer("/quota/windows/3/code"),
Some(&json!("unlimited"))
);
assert_eq!(
payload.pointer("/quota/windows/3/used_ratio"),
Some(&json!(0.3))
);
assert_eq!(
payload.pointer("/quota/windows/3/remaining_ratio"),
Some(&json!(0.7))
);
}
#[test]
fn provider_key_status_snapshot_payload_clears_expired_codex_window_exhausted_state() {
let mut key = sample_catalog_key();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "codex",
"code": "exhausted",
"exhausted": true,
"usage_ratio": 1.0,
"updated_at": 1_700_000_000u64,
"windows": [
{
"code": "weekly",
"label": "周",
"scope": "account",
"unit": "percent",
"used_ratio": 1.0,
"remaining_ratio": 0.0,
"reset_at": 1_700_003_600u64,
"window_minutes": 10_080,
"is_exhausted": true
}
]
}
}));
let payload = provider_key_status_snapshot_payload(&key, "codex");
// 汇总标志沿用既有调度口径(窗口已到期 → 不再视为耗尽)。
assert_eq!(payload.pointer("/quota/exhausted"), Some(&json!(false)));
assert_eq!(payload.pointer("/quota/code"), Some(&json!("ok")));
// 窗口级耗尽标记与剩余比例同步归一化。
assert_eq!(
payload.pointer("/quota/windows/0/is_exhausted"),
Some(&json!(false))
);
assert_eq!(
payload.pointer("/quota/windows/0/used_ratio"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_ratio"),
Some(&json!(1.0))
);
}
#[test]
fn provider_key_status_snapshot_payload_normalizes_expired_quota_windows_for_other_providers() {
let mut key = sample_catalog_key();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "kiro",
"code": "ok",
"exhausted": false,
"observed_at": 1_700_000_000u64,
"updated_at": 1_700_000_000u64,
"windows": [
{
"code": "usage",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 0.4,
"remaining_ratio": 0.6,
"used_value": 60.0,
"remaining_value": 90.0,
"limit_value": 150.0,
"reset_at": 1_700_003_600u64,
"reset_seconds": 3_600u64
},
{
"code": "usage_null_ratio",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": null,
"remaining_ratio": null,
"used_value": 60.0,
"remaining_value": 90.0,
"limit_value": 150.0,
"reset_at": 1_700_003_600u64,
"reset_seconds": 3_600u64
},
{
"code": "usage_active",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 0.4,
"remaining_ratio": 0.6,
"used_value": 60.0,
"remaining_value": 90.0,
"limit_value": 150.0,
"reset_at": 2_000_000_000u64,
"reset_seconds": 300_000_000u64
},
{
"code": "usage_no_observation",
"label": "额度",
"scope": "account",
"unit": "count",
"reset_at": 1_700_003_600u64,
"reset_seconds": 3_600u64
},
{
"code": "usage_no_deadline",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 0.4,
"remaining_ratio": 0.6,
"used_value": 60.0,
"remaining_value": 90.0,
"limit_value": 150.0
},
{
"code": "usage_remaining_without_limit",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": null,
"remaining_ratio": null,
"remaining_value": 0.0,
"reset_at": 1_700_003_600u64,
"reset_seconds": 3_600u64
}
]
}
}));
let payload = provider_key_status_snapshot_payload(&key, "kiro");
// 已到期窗口:比例与数值一起按“已重置”归一化(不再只覆盖 Codex)。
assert_eq!(
payload.pointer("/quota/windows/0/used_ratio"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_ratio"),
Some(&json!(1.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/used_value"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_value"),
Some(&json!(150.0))
);
// 比例字段为 null 时补齐为已重置口径,保证优先读比例的展示实现一致。
assert_eq!(
payload.pointer("/quota/windows/1/used_ratio"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/1/remaining_ratio"),
Some(&json!(1.0))
);
assert_eq!(
payload.pointer("/quota/windows/1/used_value"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/1/remaining_value"),
Some(&json!(150.0))
);
// 未到期窗口保留原始观测值。
assert_eq!(
payload.pointer("/quota/windows/2/used_ratio"),
Some(&json!(0.4))
);
assert_eq!(
payload.pointer("/quota/windows/2/remaining_value"),
Some(&json!(90.0))
);
// 没有用量观测的窗口不会被凭空补成 100%。
assert_eq!(payload.pointer("/quota/windows/3/used_ratio"), None);
assert_eq!(payload.pointer("/quota/windows/3/remaining_ratio"), None);
// 没有重置时间的窗口无法判定是否已重置,保持原样。
assert_eq!(
payload.pointer("/quota/windows/4/used_ratio"),
Some(&json!(0.4))
);
assert_eq!(
payload.pointer("/quota/windows/4/remaining_ratio"),
Some(&json!(0.6))
);
// 只有 remaining_value 却没有上限的窗口观测不完整,跳过归一化,
// 避免比例已恢复 100% 而数值仍停留在旧值。
assert_eq!(
payload.pointer("/quota/windows/5/remaining_ratio"),
Some(&json!(null))
);
assert_eq!(
payload.pointer("/quota/windows/5/remaining_value"),
Some(&json!(0.0))
);
}
#[test]
fn provider_key_status_snapshot_payload_clears_expired_quota_exhausted_state_for_other_providers(
) {
let mut key = sample_catalog_key();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "kiro",
"code": "exhausted",
"label": "额度耗尽",
"reason": "额度已耗尽",
"exhausted": true,
"usage_ratio": 1.0,
"observed_at": 1_700_000_000u64,
"updated_at": 1_700_000_000u64,
"windows": [
{
"code": "usage",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 1.0,
"remaining_ratio": 0.0,
"used_value": 150.0,
"remaining_value": 0.0,
"limit_value": 150.0,
"reset_at": 1_700_003_600u64,
"reset_seconds": 3_600u64,
"is_exhausted": true
}
]
}
}));
let payload = provider_key_status_snapshot_payload(&key, "kiro");
// 汇总状态沿用调度口径(窗口已到期 → 不再视为耗尽),对所有提供商生效。
assert_eq!(payload.pointer("/quota/exhausted"), Some(&json!(false)));
assert_eq!(payload.pointer("/quota/code"), Some(&json!("ok")));
// 窗口比例、数值与耗尽标记同步归一化。
assert_eq!(
payload.pointer("/quota/windows/0/is_exhausted"),
Some(&json!(false))
);
assert_eq!(
payload.pointer("/quota/windows/0/used_ratio"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_ratio"),
Some(&json!(1.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/used_value"),
Some(&json!(0.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_value"),
Some(&json!(150.0))
);
}
#[test]
fn provider_key_status_snapshot_payload_keeps_active_quota_exhausted_state_for_other_providers()
{
let mut key = sample_catalog_key();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "kiro",
"code": "exhausted",
"label": "额度耗尽",
"exhausted": true,
"usage_ratio": 1.0,
"observed_at": 1_700_000_000u64,
"updated_at": 1_700_000_000u64,
"windows": [
{
"code": "usage",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 1.0,
"remaining_ratio": 0.0,
"used_value": 150.0,
"remaining_value": 0.0,
"limit_value": 150.0,
"reset_at": 2_000_000_000u64,
"reset_seconds": 300_000_000u64,
"is_exhausted": true
}
]
}
}));
let payload = provider_key_status_snapshot_payload(&key, "kiro");
// 未到期的耗尽状态保留,不能被读取层提前“恢复”。
assert_eq!(payload.pointer("/quota/exhausted"), Some(&json!(true)));
assert_eq!(payload.pointer("/quota/code"), Some(&json!("exhausted")));
assert_eq!(
payload.pointer("/quota/windows/0/is_exhausted"),
Some(&json!(true))
);
assert_eq!(
payload.pointer("/quota/windows/0/used_ratio"),
Some(&json!(1.0))
);
assert_eq!(
payload.pointer("/quota/windows/0/remaining_ratio"),
Some(&json!(0.0))
);
}
#[test]
fn provider_key_status_snapshot_payload_backfills_chatgpt_web_image_quota() {
let mut key = sample_catalog_key();
@@ -3911,14 +4406,14 @@ mod tests {
"used_percent": 60.0,
"remaining": 60.0,
"total": 150.0,
"reset_at": 1_778_157_172u64,
"reset_at": 2_000_000_000u64,
"is_exhausted": false
},
"quota_heavy": {
"display_name": "heavy",
"remaining_fraction": 0.0,
"used_percent": 100.0,
"reset_at": 1_778_157_172u64,
"reset_at": 2_000_000_000u64,
"is_exhausted": true
}
}
@@ -3941,7 +4436,7 @@ mod tests {
assert_eq!(quota.get("pool_tier"), Some(&json!("heavy")));
assert_eq!(quota.get("exhausted"), Some(&json!(false)));
assert_eq!(quota.get("usage_ratio"), Some(&json!(1.0)));
assert_eq!(quota.get("reset_at"), Some(&json!(1_778_157_172u64)));
assert_eq!(quota.get("reset_at"), Some(&json!(2_000_000_000u64)));
assert_eq!(windows.len(), 2);
assert!(windows.iter().any(|window| {
window
@@ -4054,8 +4549,8 @@ mod tests {
"plan_name": "Pro",
"daily_remaining_percent": 40.0,
"weekly_remaining_percent": 65.0,
"daily_reset_at": 1_778_100_000u64,
"weekly_reset_at": 1_778_600_000u64,
"daily_reset_at": 2_000_000_000u64,
"weekly_reset_at": 2_000_600_000u64,
"prompt_used": 12.0,
"prompt_limit": 100.0,
"prompt_remaining": 88.0,
@@ -4090,13 +4585,13 @@ mod tests {
assert_eq!(quota.get("code"), Some(&json!("ok")));
assert_eq!(quota.get("plan_type"), Some(&json!("Pro")));
assert_eq!(quota.get("usage_ratio"), Some(&json!(0.6)));
assert_eq!(quota.get("reset_at"), Some(&json!(1_778_100_000u64)));
assert_eq!(quota.get("reset_at"), Some(&json!(2_000_000_000u64)));
assert_eq!(daily.get("remaining_ratio"), Some(&json!(0.4)));
assert_eq!(daily.get("used_ratio"), Some(&json!(0.6)));
assert_eq!(daily.get("reset_seconds"), Some(&json!(32_754u64)));
assert_eq!(daily.get("reset_seconds"), Some(&json!(221_932_754u64)));
assert_eq!(weekly.get("remaining_ratio"), Some(&json!(0.65)));
assert_eq!(weekly.get("used_ratio"), Some(&json!(0.35)));
assert_eq!(weekly.get("reset_seconds"), Some(&json!(532_754u64)));
assert_eq!(weekly.get("reset_seconds"), Some(&json!(222_532_754u64)));
assert_eq!(quota.get("allowed_models_count"), Some(&json!(82)));
}
+1
View File
@@ -92,6 +92,7 @@ mod upstream_admission;
mod usage;
mod video_tasks;
mod wallet_runtime;
mod xai_profile;
pub use self::ai_serving::api::{codex_client_originator, codex_client_user_agent};
pub(crate) use self::ai_serving::api::{
+14
View File
@@ -2527,6 +2527,20 @@ async fn run() -> Result<(), Box<dyn std::error::Error>> {
);
}
}
match state.prewarm_xai_client_profile().await {
Ok(version) => {
info!(
xai_client_version = %version,
"prewarmed Grok CLI client profile"
);
}
Err(err) => {
warn!(
error = %err,
"failed to refresh Grok CLI client profile; built-in or cached profile remains active"
);
}
}
match prewarm_direct_h2c_sender_cache_from_env_for_startup().await {
Ok(Some(report)) => {
if report.failed_targets > 0 {
+9
View File
@@ -76,6 +76,7 @@ use crate::maintenance::spawn_stats_hourly_aggregation_worker;
use crate::maintenance::spawn_usage_cleanup_worker;
use crate::maintenance::spawn_usage_counter_flush_worker;
use crate::maintenance::spawn_wallet_daily_usage_aggregation_worker;
use crate::xai_profile::spawn_worker as spawn_xai_client_profile_worker;
const SYSTEM_CONFIG_CACHE_TTL: Duration = Duration::from_secs(30);
// Requests may use a stale value after the fresh window until the entry reaches
@@ -154,6 +155,10 @@ impl AppState {
crate::codex_profile::prewarm(self.runtime_state()).await
}
pub async fn prewarm_xai_client_profile(&self) -> Result<String, String> {
crate::xai_profile::prewarm(self.runtime_state()).await
}
pub async fn prewarm_chat_pii_redaction_runtime_config(&self) -> Result<bool, String> {
crate::privacy::read_chat_pii_redaction_runtime_config(self)
.await
@@ -2356,6 +2361,10 @@ impl AppState {
crate::task_runtime::TASK_KEY_CODEX_CLIENT_PROFILE,
Some(spawn_codex_client_profile_worker(background_state.clone())),
);
supervise_worker(
crate::task_runtime::TASK_KEY_XAI_CLIENT_PROFILE,
Some(spawn_xai_client_profile_worker(background_state.clone())),
);
supervise_worker(
crate::task_runtime::TASK_KEY_VIDEO_TASK_POLLER,
spawn_video_task_poller(background_state.clone()),
@@ -25,6 +25,7 @@ pub(crate) const TASK_KEY_USAGE_COUNTER_FLUSH: &str = "usage.counter.flush.worke
pub(crate) const TASK_KEY_VIDEO_TASK_POLLER: &str = "video.task.poller";
pub(crate) const TASK_KEY_MODEL_FETCH_WORKER: &str = "model.fetch.worker";
pub(crate) const TASK_KEY_CODEX_CLIENT_PROFILE: &str = "maintenance.codex.client.profile";
pub(crate) const TASK_KEY_XAI_CLIENT_PROFILE: &str = "maintenance.xai.client.profile";
pub(crate) const TASK_KEY_PROVIDER_QUOTA_RESET: &str = "provider.quota.reset.worker";
pub(crate) const TASK_KEY_ACCOUNT_SELF_CHECK: &str = "account.self_check.worker";
pub(crate) const TASK_KEY_POOL_SCORE_REBUILD: &str = "pool.score.rebuild.worker";
@@ -211,6 +212,14 @@ const TASK_DEFINITIONS: &[TaskDefinition] = &[
true,
RETRY_ONCE,
),
TaskDefinition::new(
TASK_KEY_XAI_CLIENT_PROFILE,
TaskKind::Scheduled,
"interval",
true,
true,
RETRY_ONCE,
),
TaskDefinition::new(
TASK_KEY_PROVIDER_QUOTA_RESET,
TaskKind::Scheduled,
@@ -3044,7 +3044,94 @@ async fn gateway_prefers_status_snapshot_kiro_quota_over_stale_metadata() {
assert_eq!(keys[0]["scheduling_status"], json!("available"));
assert_eq!(keys[0]["scheduling_reason"], json!("available"));
assert_eq!(keys[0]["quota_updated_at"], json!(1_775_553_285u64));
assert_eq!(keys[0]["account_quota"], json!("剩余 75.0% (5/20)"));
// 窗口重置时间已过:读取层按“已重置”口径归一化(调度侧同样不再视为耗尽),
// 同时仍优先使用快照(上限 20)而不是陈旧元数据(上限 100)。
assert_eq!(keys[0]["account_quota"], json!("剩余 100.0% (0/20)"));
}
#[tokio::test]
async fn gateway_kiro_pool_filter_clears_stale_exhausted_summary_after_window_reset() {
let mut provider = sample_provider("provider-kiro", "kiro", 10);
provider.provider_type = "kiro".to_string();
provider.config = Some(json!({"pool_advanced": {
"reserve_minimum_quota": false, "skip_exhausted_accounts": true
}}));
let mut key = sample_key(
"key-kiro-expired",
"provider-kiro",
"claude:messages",
"oauth-placeholder",
);
key.name = "kiro expired quota key".to_string();
key.auth_type = "oauth".to_string();
key.status_snapshot = Some(json!({
"quota": {
"version": 2,
"provider_type": "kiro",
"code": "exhausted",
"label": "额度耗尽",
"reason": "额度已耗尽",
"freshness": "fresh",
"source": "refresh_api",
"observed_at": 1_775_553_285u64,
"exhausted": true,
"usage_ratio": 1.0,
"updated_at": 1_775_553_285u64,
"reset_seconds": 0u64,
"plan_type": "KIRO PRO+",
"windows": [
{
"code": "usage",
"label": "额度",
"scope": "account",
"unit": "count",
"used_ratio": 1.0,
"remaining_ratio": 0.0,
"used_value": 20.0,
"remaining_value": 0.0,
"limit_value": 20.0,
"reset_at": 1_775_639_685u64,
"reset_seconds": 0u64,
"is_exhausted": true
}
]
}
}));
let state = AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(GatewayDataState::with_provider_catalog_reader_for_tests(
Arc::new(InMemoryProviderCatalogReadRepository::seed(
vec![provider],
Vec::new(),
vec![key],
)),
));
for status in ["all", "available", "quota_exhausted"] {
let response = local_admin_pool_response(
&state,
http::Method::GET,
&format!("/api/admin/pool/provider-kiro/keys?page=1&page_size=50&status={status}"),
None,
)
.await;
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = serde_json::from_slice(
&to_bytes(response.into_body(), usize::MAX)
.await
.expect("body should read"),
)
.expect("json body should parse");
let keys = payload["keys"].as_array().expect("keys should be array");
assert_eq!(keys.len(), usize::from(status != "quota_exhausted"));
if let Some(key) = keys.first() {
assert_eq!(key["scheduling_status"], json!("available"));
assert_eq!(key["status_snapshot"]["quota"]["code"], json!("ok"));
assert_eq!(key["status_snapshot"]["quota"]["exhausted"], json!(false));
assert_eq!(key["account_quota"], json!("剩余 100.0% (0/20)"));
}
}
}
#[tokio::test]
+592
View File
@@ -0,0 +1,592 @@
//! Grok CLI 客户端版本的运行时发布与官方版本刷新。
//!
//! cli-chat-proxy.grok.com 会对低于最低版本的 `x-grok-client-version` 直接返回 426,
//! 因此网关定期读取官方发布渠道并原子替换传输层使用的版本号。
use std::collections::BTreeMap;
use std::future::Future;
use std::time::Duration;
use aether_runtime_state::RuntimeState;
use futures_util::StreamExt as _;
use reqwest::{redirect::Policy, Client};
use semver::Version;
use serde::{Deserialize, Serialize};
use tracing::{info, warn};
use crate::provider_transport::{set_xai_client_version, xai_client_version};
use crate::AppState;
/// 官方安装脚本读取的 stable 渠道,响应体是纯文本版本号。
const CLI_STABLE_CHANNEL_ENDPOINT: &str = "https://x.ai/cli/stable";
/// stable 渠道不可达时(部分部署地区无法直连 x.ai)退回 npm 发布元数据。
const CLI_NPM_RELEASE_ENDPOINT: &str = "https://registry.npmjs.org/@xai-official%2Fgrok/latest";
const CLI_NPM_PACKAGE: &str = "@xai-official/grok";
const PROFILE_CACHE_KEY: &str = "aether:xai:client-profile:v1";
const PROFILE_CACHE_TTL: Duration = Duration::from_secs(30 * 24 * 60 * 60);
/// xAI 会在发布后很快抬高最低版本,刷新间隔比 Codex 更短。
const PROFILE_REFRESH_INTERVAL: Duration = Duration::from_secs(3 * 60 * 60);
const RELEASE_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
const RELEASE_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
const MAX_RELEASE_BYTES: usize = 256 * 1024;
const CLI_TARGETS: [&str; 6] = [
"darwin-arm64",
"darwin-x64",
"linux-arm64",
"linux-x64",
"win32-arm64",
"win32-x64",
];
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct NpmRelease {
name: String,
version: String,
optional_dependencies: BTreeMap<String, String>,
}
#[derive(Debug, Deserialize, Serialize)]
struct CachedProfile {
version: String,
verified_at_unix_secs: u64,
}
#[derive(Debug, thiserror::Error)]
enum ProfileRefreshError {
#[error("Grok CLI release client initialization failed: {0}")]
Client(#[from] reqwest::Error),
#[error("Grok CLI release request returned HTTP {0}")]
HttpStatus(u16),
#[error("Grok CLI release response exceeded {MAX_RELEASE_BYTES} bytes")]
ResponseTooLarge,
#[error("Grok CLI release metadata is invalid")]
InvalidMetadata,
#[error("Grok CLI release version is older than the active profile")]
Rollback,
#[error("Grok CLI profile cache operation failed: {0}")]
Cache(String),
#[error("Grok CLI stable channel failed ({stable}); npm fallback failed ({npm})")]
AllSourcesFailed { stable: String, npm: String },
}
fn version_sequence(version: &str) -> Result<u64, ProfileRefreshError> {
let parsed = Version::parse(version).map_err(|_| ProfileRefreshError::InvalidMetadata)?;
if !parsed.pre.is_empty()
|| !parsed.build.is_empty()
|| parsed.major > 999
|| parsed.minor > 999
|| parsed.patch > 999
{
return Err(ProfileRefreshError::InvalidMetadata);
}
Ok(1 + parsed.major * 1_000_000 + parsed.minor * 1_000 + parsed.patch)
}
/// stable 渠道只返回一行版本号;任何多余内容都视为异常响应(例如被劫持的 HTML 页面)。
fn parse_stable_channel(bytes: &[u8]) -> Result<String, ProfileRefreshError> {
if bytes.len() > MAX_RELEASE_BYTES {
return Err(ProfileRefreshError::ResponseTooLarge);
}
let text = std::str::from_utf8(bytes).map_err(|_| ProfileRefreshError::InvalidMetadata)?;
let version = text.trim();
version_sequence(version)?;
Ok(version.to_owned())
}
/// 校验 npm latest 标签及六个平台二进制包来自同一版本发布。
fn parse_npm_release(bytes: &[u8]) -> Result<String, ProfileRefreshError> {
if bytes.len() > MAX_RELEASE_BYTES {
return Err(ProfileRefreshError::ResponseTooLarge);
}
let release = serde_json::from_slice::<NpmRelease>(bytes)
.map_err(|_| ProfileRefreshError::InvalidMetadata)?;
version_sequence(&release.version)?;
if release.name != CLI_NPM_PACKAGE
|| CLI_TARGETS.iter().any(|target| {
release
.optional_dependencies
.get(&format!("{CLI_NPM_PACKAGE}-{target}"))
!= Some(&release.version)
})
{
return Err(ProfileRefreshError::InvalidMetadata);
}
Ok(release.version)
}
fn refresh_enabled_from(value: Option<&str>) -> bool {
!value.is_some_and(|value| {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"0" | "false" | "off"
)
})
}
fn refresh_enabled() -> bool {
refresh_enabled_from(
std::env::var("AETHER_XAI_CLIENT_PROFILE_REFRESH")
.ok()
.as_deref(),
)
}
fn fixed_version_from(value: Option<&str>) -> Option<String> {
let value = value?.trim();
if value.is_empty() || version_sequence(value).is_err() {
None
} else {
Some(value.to_owned())
}
}
fn fixed_version_override() -> Option<String> {
let value = std::env::var("AETHER_XAI_CLIENT_VERSION").ok()?;
let version = fixed_version_from(Some(&value));
if version.is_none() {
warn!(
event_name = "xai_client_profile_fixed_version_invalid",
"AETHER_XAI_CLIENT_VERSION is invalid; using cached or built-in profile"
);
}
version
}
fn build_release_client() -> Result<Client, ProfileRefreshError> {
Client::builder()
.https_only(true)
.no_proxy()
.redirect(Policy::none())
.connect_timeout(RELEASE_CONNECT_TIMEOUT)
.timeout(RELEASE_REQUEST_TIMEOUT)
.build()
.map_err(ProfileRefreshError::Client)
}
async fn fetch_bounded(client: &Client, url: &str) -> Result<Vec<u8>, ProfileRefreshError> {
let response = client
.get(url)
.send()
.await
.map_err(ProfileRefreshError::Client)?;
if !response.status().is_success() {
return Err(ProfileRefreshError::HttpStatus(response.status().as_u16()));
}
if response
.content_length()
.is_some_and(|length| length > MAX_RELEASE_BYTES as u64)
{
return Err(ProfileRefreshError::ResponseTooLarge);
}
let mut bytes = Vec::new();
let mut stream = response.bytes_stream();
while let Some(chunk) = stream.next().await {
let chunk = chunk.map_err(ProfileRefreshError::Client)?;
if bytes.len().saturating_add(chunk.len()) > MAX_RELEASE_BYTES {
return Err(ProfileRefreshError::ResponseTooLarge);
}
bytes.extend_from_slice(&chunk);
}
Ok(bytes)
}
async fn fetch_latest_with_fallback<S, SFut, N, NFut>(
fetch_stable: S,
fetch_npm: N,
) -> Result<String, ProfileRefreshError>
where
S: FnOnce() -> SFut,
SFut: Future<Output = Result<String, ProfileRefreshError>>,
N: FnOnce() -> NFut,
NFut: Future<Output = Result<String, ProfileRefreshError>>,
{
let stable_error = match fetch_stable().await {
Ok(version) => return Ok(version),
Err(error) => error,
};
fetch_npm()
.await
.map_err(|npm_error| ProfileRefreshError::AllSourcesFailed {
stable: stable_error.to_string(),
npm: npm_error.to_string(),
})
}
async fn fetch_latest_cli_version(client: &Client) -> Result<String, ProfileRefreshError> {
fetch_latest_with_fallback(
|| async {
let bytes = fetch_bounded(client, CLI_STABLE_CHANNEL_ENDPOINT).await?;
parse_stable_channel(&bytes)
},
|| async {
let bytes = fetch_bounded(client, CLI_NPM_RELEASE_ENDPOINT).await?;
parse_npm_release(&bytes)
},
)
.await
}
fn publish_version(version: &str) -> Result<(), ProfileRefreshError> {
set_xai_client_version(version)
.map(|_| ())
.map_err(|_| ProfileRefreshError::InvalidMetadata)
}
async fn restore_cached_profile(runtime: &RuntimeState) -> Result<(), ProfileRefreshError> {
let Some(raw) = runtime
.kv_get(PROFILE_CACHE_KEY)
.await
.map_err(|err| ProfileRefreshError::Cache(err.to_string()))?
else {
return Ok(());
};
let cached = serde_json::from_str::<CachedProfile>(&raw)
.map_err(|_| ProfileRefreshError::InvalidMetadata)?;
if let Some(version) = cached_version_to_restore(&cached, &xai_client_version())? {
publish_version(&version)?;
info!(
event_name = "xai_client_profile_restored",
version = %version,
verified_at_unix_secs = cached.verified_at_unix_secs,
"restored cached Grok CLI profile"
);
}
Ok(())
}
fn cached_version_to_restore(
cached: &CachedProfile,
active_version: &str,
) -> Result<Option<String>, ProfileRefreshError> {
let cached_sequence = version_sequence(&cached.version)?;
let active_sequence = version_sequence(active_version)?;
Ok((cached_sequence > active_sequence).then(|| cached.version.clone()))
}
async fn refresh_once_with_fetch<F, Fut>(
runtime: &RuntimeState,
fixed_version: Option<&str>,
refresh_is_enabled: bool,
fetch_latest: F,
) -> Result<String, ProfileRefreshError>
where
F: FnOnce() -> Fut,
Fut: Future<Output = Result<String, ProfileRefreshError>>,
{
if let Some(version) = fixed_version {
publish_version(version)?;
return Ok(version.to_owned());
}
if let Err(error) = restore_cached_profile(runtime).await {
// 缓存损坏或暂时不可用不应阻断官方版本检查;当前进程继续使用旧画像。
warn!(
event_name = "xai_client_profile_cache_restore_failed",
error = %error,
"could not restore cached Grok CLI profile"
);
}
if !refresh_is_enabled {
return Ok(xai_client_version());
}
let version = fetch_latest().await?;
let current = xai_client_version();
if version_sequence(&version)? < version_sequence(&current)? {
return Err(ProfileRefreshError::Rollback);
}
let cached = CachedProfile {
version: version.clone(),
verified_at_unix_secs: chrono::Utc::now().timestamp().max(0) as u64,
};
let serialized =
serde_json::to_string(&cached).map_err(|_| ProfileRefreshError::InvalidMetadata)?;
publish_version(&version)?;
if let Err(error) = runtime
.kv_set(PROFILE_CACHE_KEY, serialized, Some(PROFILE_CACHE_TTL))
.await
{
// 本地版本已经完成原子替换;缓存写失败只影响下次进程启动的恢复。
warn!(
event_name = "xai_client_profile_cache_write_failed",
error = %error,
"published Grok CLI profile locally but could not persist the cache"
);
}
Ok(version)
}
async fn refresh_once(runtime: &RuntimeState) -> Result<String, ProfileRefreshError> {
let fixed_version = fixed_version_override();
refresh_once_with_fetch(
runtime,
fixed_version.as_deref(),
refresh_enabled(),
|| async {
let client = build_release_client()?;
fetch_latest_cli_version(&client).await
},
)
.await
}
pub(crate) async fn prewarm(runtime: &RuntimeState) -> Result<String, String> {
refresh_once(runtime).await.map_err(|err| err.to_string())
}
pub(crate) fn spawn_worker(app: AppState) -> tokio::task::JoinHandle<()> {
crate::task_runtime::spawn_singleton_worker(
app,
crate::task_runtime::TASK_KEY_XAI_CLIENT_PROFILE,
|app| async move {
let mut interval = tokio::time::interval(PROFILE_REFRESH_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
// 启动阶段由 prewarm 完成一次检查;后台任务只负责后续定时刷新,避免重复建连。
interval.tick().await;
loop {
interval.tick().await;
match refresh_once(app.runtime_state()).await {
Ok(version) => info!(
event_name = "xai_client_profile_refreshed",
version = %version,
"refreshed Grok CLI profile"
),
Err(error) => warn!(
event_name = "xai_client_profile_refresh_failed",
error = %error,
"keeping the previous Grok CLI profile after refresh failure"
),
}
}
},
)
}
#[cfg(test)]
mod tests {
use std::sync::{
atomic::{AtomicBool, Ordering},
Mutex, OnceLock,
};
use std::time::Duration;
use aether_runtime_state::{MemoryRuntimeStateConfig, RuntimeState};
use super::{
cached_version_to_restore, fetch_latest_with_fallback, fixed_version_from,
parse_npm_release, parse_stable_channel, refresh_enabled_from, refresh_once_with_fetch,
CachedProfile, ProfileRefreshError, PROFILE_CACHE_KEY,
};
use crate::provider_transport::{set_xai_client_version, xai_client_version};
static PROFILE_TEST_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
struct VersionRestore(String);
impl Drop for VersionRestore {
fn drop(&mut self) {
let _ = set_xai_client_version(&self.0);
}
}
fn version_restore_guard() -> (std::sync::MutexGuard<'static, ()>, VersionRestore) {
let lock = PROFILE_TEST_LOCK.get_or_init(|| Mutex::new(()));
let guard = lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let restore = VersionRestore(xai_client_version());
(guard, restore)
}
fn npm_release(version: &str) -> serde_json::Value {
let mut deps = serde_json::Map::new();
for target in super::CLI_TARGETS {
deps.insert(
format!("@xai-official/grok-{target}"),
serde_json::Value::String(version.to_string()),
);
}
serde_json::json!({
"name": "@xai-official/grok",
"version": version,
"optionalDependencies": deps,
})
}
#[test]
fn stable_channel_accepts_only_a_bare_release_version() {
assert_eq!(parse_stable_channel(b"1.0.46\n").unwrap(), "1.0.46");
assert!(parse_stable_channel(b"<html>1.0.46</html>").is_err());
assert!(parse_stable_channel(b"1.0.47-alpha.1").is_err());
assert!(parse_stable_channel(b"").is_err());
}
#[test]
fn npm_release_requires_every_platform_binary_at_the_same_version() {
let body = npm_release("1.0.46");
assert_eq!(
parse_npm_release(&serde_json::to_vec(&body).unwrap()).unwrap(),
"1.0.46"
);
let mut mismatched = npm_release("1.0.46");
mismatched["optionalDependencies"]["@xai-official/grok-linux-x64"] =
serde_json::Value::String("1.0.45".to_string());
assert!(parse_npm_release(&serde_json::to_vec(&mismatched).unwrap()).is_err());
let mut wrong_package = npm_release("1.0.46");
wrong_package["name"] = serde_json::Value::String("grok".to_string());
assert!(parse_npm_release(&serde_json::to_vec(&wrong_package).unwrap()).is_err());
}
#[test]
fn refresh_and_fixed_version_environment_policies_are_strict() {
assert!(!refresh_enabled_from(Some("off")));
assert!(!refresh_enabled_from(Some(" FALSE ")));
assert!(refresh_enabled_from(None));
assert_eq!(
fixed_version_from(Some(" 1.0.46 ")).as_deref(),
Some("1.0.46")
);
assert!(fixed_version_from(Some("1.0.46-beta.1")).is_none());
assert!(fixed_version_from(Some("1.0")).is_none());
}
#[test]
fn cached_profile_never_rewinds_active_profile() {
let cached = CachedProfile {
version: "1.0.50".to_string(),
verified_at_unix_secs: 1,
};
assert_eq!(
cached_version_to_restore(&cached, "1.0.46").unwrap(),
Some("1.0.50".to_string())
);
assert_eq!(cached_version_to_restore(&cached, "1.1.0").unwrap(), None);
}
#[tokio::test]
async fn npm_is_used_only_when_the_stable_channel_fails() {
let npm_called = AtomicBool::new(false);
let version = fetch_latest_with_fallback(
|| async { Ok("1.0.46".to_string()) },
|| async {
npm_called.store(true, Ordering::SeqCst);
Ok("1.0.45".to_string())
},
)
.await
.unwrap();
assert_eq!(version, "1.0.46");
assert!(!npm_called.load(Ordering::SeqCst));
let version = fetch_latest_with_fallback(
|| async { Err(ProfileRefreshError::HttpStatus(503)) },
|| async { Ok("1.0.46".to_string()) },
)
.await
.unwrap();
assert_eq!(version, "1.0.46");
let result = fetch_latest_with_fallback(
|| async { Err(ProfileRefreshError::HttpStatus(503)) },
|| async { Err(ProfileRefreshError::InvalidMetadata) },
)
.await;
assert!(matches!(
result,
Err(ProfileRefreshError::AllSourcesFailed { .. })
));
}
#[tokio::test]
async fn cache_hit_is_restored_without_network_when_refresh_is_disabled() {
let (_lock, _restore) = version_restore_guard();
set_xai_client_version("1.0.46").unwrap();
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
runtime
.kv_set(
PROFILE_CACHE_KEY,
serde_json::to_string(&CachedProfile {
version: "1.0.50".to_string(),
verified_at_unix_secs: 1,
})
.unwrap(),
Some(Duration::from_secs(60)),
)
.await
.unwrap();
let result = refresh_once_with_fetch(&runtime, None, false, || async {
Err(ProfileRefreshError::HttpStatus(599))
})
.await
.unwrap();
assert_eq!(result, "1.0.50");
assert_eq!(xai_client_version(), "1.0.50");
}
#[tokio::test]
async fn refresh_failure_keeps_previous_profile() {
let (_lock, _restore) = version_restore_guard();
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let before = xai_client_version();
let result = refresh_once_with_fetch(&runtime, None, true, || async {
Err(ProfileRefreshError::HttpStatus(503))
})
.await;
assert!(matches!(result, Err(ProfileRefreshError::HttpStatus(503))));
assert_eq!(xai_client_version(), before);
}
#[tokio::test]
async fn successful_refresh_publishes_and_caches_version() {
let (_lock, _restore) = version_restore_guard();
set_xai_client_version("1.0.46").unwrap();
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let result =
refresh_once_with_fetch(&runtime, None, true, || async { Ok("1.0.51".to_string()) })
.await
.unwrap();
assert_eq!(result, "1.0.51");
assert_eq!(xai_client_version(), "1.0.51");
let cached = runtime.kv_get(PROFILE_CACHE_KEY).await.unwrap().unwrap();
assert!(cached.contains("\"1.0.51\""));
}
#[tokio::test]
async fn fixed_version_override_skips_network_and_publishes_version() {
let (_lock, _restore) = version_restore_guard();
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let fetch_called = AtomicBool::new(false);
let result = refresh_once_with_fetch(&runtime, Some("1.0.60"), true, || async {
fetch_called.store(true, Ordering::SeqCst);
Ok("1.0.61".to_string())
})
.await
.unwrap();
assert_eq!(result, "1.0.60");
assert!(!fetch_called.load(Ordering::SeqCst));
assert_eq!(xai_client_version(), "1.0.60");
}
#[tokio::test]
async fn rollback_is_rejected_without_replacing_profile() {
let (_lock, _restore) = version_restore_guard();
set_xai_client_version("1.0.60").unwrap();
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let result =
refresh_once_with_fetch(&runtime, None, true, || async { Ok("1.0.59".to_string()) })
.await;
assert!(matches!(result, Err(ProfileRefreshError::Rollback)));
assert_eq!(xai_client_version(), "1.0.60");
}
}