refactor: unify provider client identity profiles and node synchronization

This commit is contained in:
dalamudx
2026-10-05 15:18:36 +08:00
parent 716c35bf56
commit 7acbaab82b
35 changed files with 1006 additions and 773 deletions
@@ -7,7 +7,7 @@ use crate::ai_serving::transport::{
build_gemini_cli_v1internal_request, build_standard_provider_request_headers,
GatewayProviderTransportSnapshot, GeminiCliRequestAuth, GeminiCliRequestAuthSupport,
GeminiCliRequestEnvelopeSupport, StandardProviderRequestHeaders,
StandardProviderRequestHeadersInput, GEMINI_CLI_USER_AGENT,
StandardProviderRequestHeadersInput,
};
use crate::AppState;
@@ -64,8 +64,10 @@ pub(crate) async fn build_gemini_cli_v1internal_provider_request(
)
.ok_or(GeminiCliV1InternalRequestError::UpstreamUrlUnavailable)?;
let extra_headers =
BTreeMap::from([("user-agent".to_string(), GEMINI_CLI_USER_AGENT.to_string())]);
let extra_headers = BTreeMap::from([(
"user-agent".to_string(),
aether_provider_transport::gemini_cli::gemini_cli_client_user_agent(),
)]);
let headers = build_standard_provider_request_headers(StandardProviderRequestHeadersInput {
transport: &payload.transport,
provider_api_format: input.provider_api_format,
@@ -21,8 +21,7 @@ use crate::ai_serving::transport::{
build_same_format_provider_headers, resolve_local_gemini_cli_request_auth,
GeminiCliRequestAuth, GeminiCliRequestAuthSupport, GeminiCliRequestEnvelopeSupport,
GrokHeaderInput, SameFormatProviderCompatibilityEdit,
SameFormatProviderCompatibilityEditAction, SameFormatProviderHeadersInput,
GEMINI_CLI_USER_AGENT, GROK_CHAT_PATH,
SameFormatProviderCompatibilityEditAction, SameFormatProviderHeadersInput, GROK_CHAT_PATH,
};
use crate::ai_serving::{
CandidateFailureDiagnostic, GatewayProviderTransportSnapshot, CODEX_RESPONSES_LITE_HEADER,
@@ -533,7 +532,10 @@ pub(crate) async fn resolve_local_same_format_provider_candidate_payload_parts(
.map(build_antigravity_static_identity_headers)
.unwrap_or_default();
if prepared.behavior.is_gemini_cli {
extra_headers.insert("user-agent".to_string(), GEMINI_CLI_USER_AGENT.to_string());
extra_headers.insert(
"user-agent".to_string(),
aether_provider_transport::gemini_cli::gemini_cli_client_user_agent(),
);
}
let Some(mut provider_request_headers) = (if is_grok {
build_grok_browser_headers(GrokHeaderInput {
+461 -19
View File
@@ -1,4 +1,4 @@
//! CLI 客户端画像(Codex / Claude Code)的运行时发布与官方版本刷新。
//! CLI 客户端画像的统一发布、官方版本刷新与每节点缓存同步。
//!
//! 每个客户端由一份 [`CliClientProfileSpec`] 描述:官方 npm 发布源、平台包校验规则、
//! 运行时缓存键、环境变量开关与画像发布函数。刷新逻辑本身与客户端无关。
@@ -23,16 +23,20 @@ use crate::AppState;
const PROFILE_CACHE_TTL: Duration = Duration::from_secs(30 * 24 * 60 * 60);
const PROFILE_REFRESH_INTERVAL: Duration = Duration::from_secs(24 * 60 * 60);
const PROFILE_SYNC_INTERVAL: Duration = Duration::from_secs(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;
/// 单个 CLI 客户端的发布源、校验规则、缓存与运行时画像发布方式。
#[derive(Clone, Copy)]
pub(crate) struct CliClientProfileSpec {
/// 日志中的客户端标识。
client: &'static str,
/// 官方 npm stable 标签的发布元数据地址。
release_endpoint: &'static str,
stable_channel: Option<&'static str>,
refresh_interval: Duration,
package_name: &'static str,
/// 同一发布必须同时携带的平台二进制包。
platform_targets: &'static [&'static str],
@@ -67,6 +71,8 @@ fn publish_claude_code_version(version: &str) -> Result<(), &'static str> {
pub(crate) static CODEX_CLI_PROFILE: CliClientProfileSpec = CliClientProfileSpec {
client: "codex",
stable_channel: None,
refresh_interval: PROFILE_REFRESH_INTERVAL,
release_endpoint: "https://registry.npmjs.org/@openai%2Fcodex/latest",
package_name: "@openai/codex",
platform_targets: &[
@@ -91,6 +97,8 @@ pub(crate) static CODEX_CLI_PROFILE: CliClientProfileSpec = CliClientProfileSpec
/// 指纹仍由 transport crate 中带版本号的身份模板统一维护。
pub(crate) static CLAUDE_CODE_CLI_PROFILE: CliClientProfileSpec = CliClientProfileSpec {
client: "claude_code",
stable_channel: None,
refresh_interval: PROFILE_REFRESH_INTERVAL,
release_endpoint: "https://registry.npmjs.org/@anthropic-ai%2Fclaude-code/latest",
package_name: "@anthropic-ai/claude-code",
platform_targets: &[
@@ -112,11 +120,72 @@ pub(crate) static CLAUDE_CODE_CLI_PROFILE: CliClientProfileSpec = CliClientProfi
publish_version: publish_claude_code_version,
};
fn publish_xai_version(version: &str) -> Result<(), &'static str> {
aether_provider_transport::xai::set_xai_client_version(version).map(|_| ())
}
fn publish_gemini_version(version: &str) -> Result<(), &'static str> {
aether_provider_transport::gemini_cli::set_gemini_cli_client_version(version).map(|_| ())
}
pub(crate) static XAI_CLI_PROFILE: CliClientProfileSpec = CliClientProfileSpec {
client: "xai",
release_endpoint: "https://registry.npmjs.org/@xai-official%2Fgrok/latest",
stable_channel: Some("https://x.ai/cli/stable"),
refresh_interval: Duration::from_secs(3 * 60 * 60),
package_name: "@xai-official/grok",
platform_targets: &[
"darwin-arm64",
"darwin-x64",
"linux-arm64",
"linux-x64",
"win32-arm64",
"win32-x64",
],
platform_dependency: claude_code_platform_dependency,
cache_key: "aether:xai:client-profile:v1",
refresh_env: "AETHER_XAI_CLIENT_PROFILE_REFRESH",
fixed_version_env: "AETHER_XAI_CLIENT_VERSION",
task_key: crate::task_runtime::TASK_KEY_XAI_CLIENT_PROFILE,
active_version: aether_provider_transport::xai::xai_client_version,
publish_version: publish_xai_version,
};
pub(crate) static GEMINI_CLI_PROFILE: CliClientProfileSpec = CliClientProfileSpec {
client: "gemini_cli",
release_endpoint: "https://registry.npmjs.org/@google%2Fgemini-cli/latest",
stable_channel: None,
refresh_interval: PROFILE_REFRESH_INTERVAL,
package_name: "@google/gemini-cli",
// Official JS bundle has no same-version platform packages.
platform_targets: &[],
platform_dependency: claude_code_platform_dependency,
cache_key: "aether:gemini_cli:client-profile:v1",
refresh_env: "AETHER_GEMINI_CLI_CLIENT_PROFILE_REFRESH",
fixed_version_env: "AETHER_GEMINI_CLI_CLIENT_VERSION",
task_key: crate::task_runtime::TASK_KEY_GEMINI_CLI_CLIENT_PROFILE,
active_version: aether_provider_transport::gemini_cli::gemini_cli_client_version,
publish_version: publish_gemini_version,
};
pub(crate) static CLI_PROFILES: &[&CliClientProfileSpec] = &[
&CODEX_CLI_PROFILE,
&CLAUDE_CODE_CLI_PROFILE,
&XAI_CLI_PROFILE,
&GEMINI_CLI_PROFILE,
];
impl CliClientProfileSpec {
pub(crate) fn task_key(&self) -> &'static str {
self.task_key
}
pub(crate) fn client_name(&self) -> &'static str {
self.client
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct NpmRelease {
name: String,
version: String,
#[serde(default)]
optional_dependencies: BTreeMap<String, String>,
}
@@ -140,6 +209,8 @@ enum ProfileRefreshError {
Rollback,
#[error("CLI profile cache operation failed: {0}")]
Cache(String),
#[error("stable channel failed ({stable}); npm fallback failed ({npm})")]
AllSourcesFailed { stable: String, npm: String },
}
fn version_sequence(version: &str) -> Result<u64, ProfileRefreshError> {
@@ -226,12 +297,9 @@ fn build_release_client() -> Result<Client, ProfileRefreshError> {
.map_err(ProfileRefreshError::Client)
}
async fn fetch_latest_cli_version(
spec: &CliClientProfileSpec,
client: &Client,
) -> Result<String, ProfileRefreshError> {
async fn fetch_bounded(client: &Client, url: &str) -> Result<Vec<u8>, ProfileRefreshError> {
let response = client
.get(spec.release_endpoint)
.get(url)
.send()
.await
.map_err(ProfileRefreshError::Client)?;
@@ -254,7 +322,56 @@ async fn fetch_latest_cli_version(
}
bytes.extend_from_slice(&chunk);
}
parse_cli_release(spec, &bytes)
Ok(bytes)
}
fn parse_stable_channel(bytes: &[u8]) -> Result<String, ProfileRefreshError> {
if bytes.len() > MAX_RELEASE_BYTES {
return Err(ProfileRefreshError::ResponseTooLarge);
}
let version = std::str::from_utf8(bytes)
.map_err(|_| ProfileRefreshError::InvalidMetadata)?
.trim();
version_sequence(version)?;
Ok(version.to_owned())
}
async fn fetch_latest_with_fallback<S, SF, N, NF>(
stable: S,
npm: N,
) -> Result<String, ProfileRefreshError>
where
S: FnOnce() -> SF,
SF: Future<Output = Result<String, ProfileRefreshError>>,
N: FnOnce() -> NF,
NF: Future<Output = Result<String, ProfileRefreshError>>,
{
match stable().await {
Ok(version) => Ok(version),
Err(stable) => npm()
.await
.map_err(|npm| ProfileRefreshError::AllSourcesFailed {
stable: stable.to_string(),
npm: npm.to_string(),
}),
}
}
async fn fetch_latest_cli_version(
spec: &CliClientProfileSpec,
client: &Client,
) -> Result<String, ProfileRefreshError> {
let npm =
|| async { parse_cli_release(spec, &fetch_bounded(client, spec.release_endpoint).await?) };
if let Some(url) = spec.stable_channel {
fetch_latest_with_fallback(
|| async { parse_stable_channel(&fetch_bounded(client, url).await?) },
npm,
)
.await
} else {
npm().await
}
}
fn publish(spec: &CliClientProfileSpec, version: &str) -> Result<(), ProfileRefreshError> {
@@ -293,7 +410,7 @@ fn cached_version_to_restore(
) -> 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()))
Ok((cached_sequence > active_sequence).then(|| cached.version.clone()))
}
async fn refresh_once_with_fetch<F, Fut>(
@@ -326,6 +443,8 @@ where
}
let version = fetch_latest().await?;
// Another publisher may have advanced shared state while the fetch awaited.
let _ = restore_cached_profile(spec, runtime).await;
let current = (spec.active_version)();
if version_sequence(&version)? < version_sequence(&current)? {
return Err(ProfileRefreshError::Rollback);
@@ -342,7 +461,7 @@ where
.kv_set(spec.cache_key, serialized, Some(PROFILE_CACHE_TTL))
.await
{
// 本地画像已经完成原子替换;缓存写失败只影响下次进程启动的恢复。
// 本地画像已经完成原子替换;缓存写失败会延迟其他节点同步及下次启动的恢复。
warn!(
event_name = "cli_client_profile_cache_write_failed",
client = spec.client,
@@ -358,17 +477,49 @@ async fn refresh_once(
runtime: &RuntimeState,
) -> Result<String, ProfileRefreshError> {
let fixed_version = fixed_version_override(spec);
refresh_once_with_fetch(
spec,
runtime,
fixed_version.as_deref(),
refresh_enabled(spec),
|| async {
if fixed_version.is_some() || !refresh_enabled(spec) {
return refresh_once_with_fetch(spec, runtime, fixed_version.as_deref(), false, || async {
Err(ProfileRefreshError::InvalidMetadata)
})
.await;
}
let _ = restore_cached_profile(spec, runtime).await;
let lock_key = format!("aether:client-profile:refresh:{}", spec.client);
let owner = uuid::Uuid::now_v7().to_string();
let Some(lease) = runtime
.lock_try_acquire(&lock_key, &owner, Duration::from_secs(120))
.await
.map_err(|e| ProfileRefreshError::Cache(e.to_string()))?
else {
return Ok((spec.active_version)());
};
let result = async {
// Concurrent startups must not all query the release source. A fresh,
// verified cache is sufficient; per-node sync never accesses the network.
if let Ok(Some(raw)) = runtime.kv_get(spec.cache_key).await {
if let Ok(cached) = serde_json::from_str::<CachedProfile>(&raw) {
let now = chrono::Utc::now().timestamp().max(0) as u64;
if now >= cached.verified_at_unix_secs
&& now - cached.verified_at_unix_secs < spec.refresh_interval.as_secs()
&& version_sequence(&cached.version).is_ok_and(|cached_seq| {
version_sequence(&(spec.active_version)())
.is_ok_and(|active_seq| cached_seq >= active_seq)
})
{
restore_cached_profile(spec, runtime).await?;
return Ok((spec.active_version)());
}
}
}
refresh_once_with_fetch(spec, runtime, None, true, || async {
let client = build_release_client()?;
fetch_latest_cli_version(spec, &client).await
},
)
.await
})
.await
}
.await;
let _ = runtime.lock_release(&lease).await;
result
}
pub(crate) async fn prewarm(
@@ -385,7 +536,7 @@ pub(crate) fn spawn_worker(
app: AppState,
) -> tokio::task::JoinHandle<()> {
crate::task_runtime::spawn_singleton_worker(app, spec.task_key, move |app| async move {
let mut interval = tokio::time::interval(PROFILE_REFRESH_INTERVAL);
let mut interval = tokio::time::interval(spec.refresh_interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
// 启动阶段由 prewarm 完成一次检查;后台任务只负责后续每日刷新,避免重复建连。
interval.tick().await;
@@ -409,6 +560,49 @@ pub(crate) fn spawn_worker(
})
}
/// One task per process (including frontdoor-only nodes), independent of the
/// cluster singleton release checkers. Dropping the guard aborts the task.
pub struct ClientProfileSyncGuard(tokio::task::JoinHandle<()>);
impl Drop for ClientProfileSyncGuard {
fn drop(&mut self) {
self.0.abort();
}
}
async fn sync_cached_profile(
spec: &CliClientProfileSpec,
runtime: &RuntimeState,
fixed: Option<&str>,
) -> Result<(), ProfileRefreshError> {
if let Some(version) = fixed {
publish(spec, version)
} else {
restore_cached_profile(spec, runtime).await
}
}
pub(crate) fn spawn_cache_sync(app: AppState) -> ClientProfileSyncGuard {
ClientProfileSyncGuard(aether_runtime::task::spawn_named(
"client-profile-cache-sync",
async move {
let mut interval = tokio::time::interval(PROFILE_SYNC_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
interval.tick().await;
for spec in CLI_PROFILES {
let fixed = fixed_version_override(spec);
if let Err(error) =
sync_cached_profile(spec, app.runtime_state(), fixed.as_deref()).await
{
warn!(event_name = "cli_client_profile_cache_sync_failed", client = spec.client, error = %error,
"keeping the previous local client profile");
}
}
}
},
))
}
#[cfg(test)]
mod tests {
use std::sync::{
@@ -448,6 +642,8 @@ mod tests {
static TEST_PROFILE: CliClientProfileSpec = CliClientProfileSpec {
client: "test",
stable_channel: None,
refresh_interval: super::PROFILE_REFRESH_INTERVAL,
release_endpoint: "https://registry.invalid/test/latest",
package_name: "@test/cli",
platform_targets: &["linux-x64"],
@@ -466,6 +662,252 @@ mod tests {
guard
}
#[test]
fn all_release_adapters_keep_distinct_cache_keys_and_registered_tasks() {
let mut keys = std::collections::HashSet::new();
let mut tasks = std::collections::HashSet::new();
for spec in super::CLI_PROFILES {
assert!(keys.insert(spec.cache_key));
assert!(tasks.insert(spec.task_key()));
assert!(spec
.release_endpoint
.starts_with("https://registry.npmjs.org/"));
}
assert_eq!(keys.len(), 4);
assert_eq!(
super::XAI_CLI_PROFILE.refresh_interval,
Duration::from_secs(3 * 60 * 60)
);
}
#[test]
fn gemini_bundle_accepts_only_the_official_stable_package() {
let mut body = serde_json::json!({"name":"@google/gemini-cli", "version":"0.62.0"});
assert_eq!(
parse_cli_release(
&super::GEMINI_CLI_PROFILE,
&serde_json::to_vec(&body).unwrap()
)
.unwrap(),
"0.62.0"
);
body["name"] = serde_json::json!("gemini-cli");
assert!(parse_cli_release(
&super::GEMINI_CLI_PROFILE,
&serde_json::to_vec(&body).unwrap()
)
.is_err());
body["name"] = serde_json::json!("@google/gemini-cli");
body["version"] = serde_json::json!("0.63.0-preview.1");
assert!(parse_cli_release(
&super::GEMINI_CLI_PROFILE,
&serde_json::to_vec(&body).unwrap()
)
.is_err());
}
#[test]
fn grok_preserves_stable_channel_and_all_platform_release_checks() {
assert_eq!(super::parse_stable_channel(b"1.0.46\n").unwrap(), "1.0.46");
for bytes in [
b"<html>1.0.46</html>".as_slice(),
b"1.0.47-alpha.1",
b"",
b"1.0.46+build",
] {
assert!(super::parse_stable_channel(bytes).is_err());
}
let spec = &super::XAI_CLI_PROFILE;
let dependencies = spec
.platform_targets
.iter()
.map(|target| {
(
format!("@xai-official/grok-{target}"),
serde_json::json!("1.0.46"),
)
})
.collect::<serde_json::Map<_, _>>();
let mut body = serde_json::json!({"name":"@xai-official/grok", "version":"1.0.46", "optionalDependencies":dependencies});
assert_eq!(
parse_cli_release(spec, &serde_json::to_vec(&body).unwrap()).unwrap(),
"1.0.46"
);
body["optionalDependencies"]["@xai-official/grok-linux-x64"] = serde_json::json!("1.0.45");
assert!(parse_cli_release(spec, &serde_json::to_vec(&body).unwrap()).is_err());
}
#[tokio::test]
async fn grok_uses_npm_only_after_stable_fails() {
let called = AtomicBool::new(false);
assert_eq!(
super::fetch_latest_with_fallback(
|| async { Ok("1.0.46".into()) },
|| async {
called.store(true, Ordering::SeqCst);
Ok("1.0.47".into())
}
)
.await
.unwrap(),
"1.0.46"
);
assert!(!called.load(Ordering::SeqCst));
assert_eq!(
super::fetch_latest_with_fallback(
|| async { Err(ProfileRefreshError::HttpStatus(503)) },
|| async { Ok("1.0.47".into()) }
)
.await
.unwrap(),
"1.0.47"
);
assert!(matches!(
super::fetch_latest_with_fallback(
|| async { Err(ProfileRefreshError::HttpStatus(503)) },
|| async { Err(ProfileRefreshError::HttpStatus(502)) }
)
.await,
Err(ProfileRefreshError::AllSourcesFailed { .. })
));
}
static REPLICA_VERSION: Mutex<String> = Mutex::new(String::new());
fn replica_version() -> String {
REPLICA_VERSION.lock().unwrap().clone()
}
fn publish_replica(version: &str) -> Result<(), &'static str> {
*REPLICA_VERSION.lock().unwrap() = version.into();
Ok(())
}
#[tokio::test]
async fn a_non_owner_replica_syncs_without_fetching_and_fixed_override_wins() {
let _guard = test_profile_guard().await;
publish_replica(TEST_BUILTIN_VERSION).unwrap();
let replica = CliClientProfileSpec {
active_version: replica_version,
publish_version: publish_replica,
..TEST_PROFILE
};
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
refresh_once_with_fetch(&TEST_PROFILE, &runtime, None, true, || async {
Ok("1.2.0".into())
})
.await
.unwrap();
assert_eq!(replica_version(), TEST_BUILTIN_VERSION);
super::sync_cached_profile(&replica, &runtime, None)
.await
.unwrap();
assert_eq!(replica_version(), "1.2.0");
super::sync_cached_profile(&replica, &runtime, Some("1.0.0"))
.await
.unwrap();
assert_eq!(replica_version(), "1.0.0");
assert!(runtime
.kv_get(TEST_PROFILE.cache_key)
.await
.unwrap()
.unwrap()
.contains("1.2.0"));
}
#[tokio::test]
async fn newer_shared_version_arriving_during_fetch_rejects_stale_publish() {
let _guard = test_profile_guard().await;
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let result = refresh_once_with_fetch(&TEST_PROFILE, &runtime, None, true, || async {
runtime
.kv_set(
TEST_PROFILE.cache_key,
serde_json::to_string(&CachedProfile {
version: "1.6.0".into(),
verified_at_unix_secs: 1,
})
.unwrap(),
None,
)
.await
.unwrap();
Ok("1.5.0".into())
})
.await;
assert!(matches!(result, Err(ProfileRefreshError::Rollback)));
assert_eq!(test_active_version(), "1.6.0");
assert!(runtime
.kv_get(TEST_PROFILE.cache_key)
.await
.unwrap()
.unwrap()
.contains("1.6.0"));
}
#[tokio::test]
async fn corrupt_cache_does_not_block_a_verified_refresh_or_erase_local_state() {
let _guard = test_profile_guard().await;
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
runtime
.kv_set(TEST_PROFILE.cache_key, "invalid-json", None)
.await
.unwrap();
assert!(super::sync_cached_profile(&TEST_PROFILE, &runtime, None)
.await
.is_err());
assert_eq!(test_active_version(), TEST_BUILTIN_VERSION);
assert_eq!(
refresh_once_with_fetch(&TEST_PROFILE, &runtime, None, true, || async {
Ok("1.2.0".into())
})
.await
.unwrap(),
"1.2.0"
);
}
#[tokio::test]
async fn a_fresh_verified_cache_avoids_a_startup_network_check() {
let _guard = test_profile_guard().await;
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
runtime
.kv_set(
TEST_PROFILE.cache_key,
serde_json::to_string(&CachedProfile {
version: "1.2.0".into(),
verified_at_unix_secs: chrono::Utc::now().timestamp().max(0) as u64,
})
.unwrap(),
None,
)
.await
.unwrap();
// TEST_PROFILE's URL cannot return metadata; success demonstrates no HTTP fetch.
assert_eq!(
super::refresh_once(&TEST_PROFILE, &runtime).await.unwrap(),
"1.2.0"
);
}
#[tokio::test]
async fn another_startup_holding_the_release_lock_skips_the_network() {
let _guard = test_profile_guard().await;
let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default());
let lease = runtime
.lock_try_acquire(
"aether:client-profile:refresh:test",
"other-node",
Duration::from_secs(120),
)
.await
.unwrap()
.unwrap();
assert_eq!(
super::refresh_once(&TEST_PROFILE, &runtime).await.unwrap(),
TEST_BUILTIN_VERSION
);
runtime.lock_release(&lease).await.unwrap();
}
#[test]
fn accepts_only_one_verified_codex_release_for_all_targets() {
let body = serde_json::json!({
@@ -43,12 +43,11 @@ use crate::AppState;
const CHATGPT_WEB_INTERNAL_HEADER: &str = "x-aether-chatgpt-web-image";
const CHATGPT_WEB_DEFAULT_BASE_URL: &str = "https://chatgpt.com";
const CHATGPT_WEB_USER_AGENT: &str = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/143.0.0.0 Safari/537.36 Edg/143.0.0.0";
const CHATGPT_WEB_CLIENT_VERSION: &str = "prod-be885abbfcfe7b1f511e88b3003d9ee44757fbad";
const CHATGPT_WEB_BUILD_NUMBER: &str = "5955942";
const CHATGPT_WEB_SEC_CH_UA: &str =
r#""Microsoft Edge";v="143", "Chromium";v="143", "Not A(Brand";v="24""#;
const CHATGPT_WEB_BROWSER_PROFILE: &str = "chrome143";
use aether_provider_transport::client_identity::CHATGPT_WEB_BROWSER_PROFILE;
use aether_provider_transport::client_identity::{
CHATGPT_WEB_BUILD_NUMBER, CHATGPT_WEB_CLIENT_VERSION, CHATGPT_WEB_SEC_CH_UA,
CHATGPT_WEB_USER_AGENT,
};
const CHATGPT_WEB_QUOTA_REFRESH_TIMEOUT_MS: u64 = 30_000;
const CHATGPT_WEB_QUOTA_REFRESH_PROXY_TIMEOUT_MS: u64 = 60_000;
const RUNTIME_METADATA_CAS_MAX_ATTEMPTS: usize = 16;
@@ -3462,7 +3462,7 @@ async fn provider_query_execute_standard_test_candidate(
{
request_headers
.entry("user-agent".to_string())
.or_insert_with(|| crate::provider_transport::GEMINI_CLI_USER_AGENT.to_string());
.or_insert_with(aether_provider_transport::gemini_cli::gemini_cli_client_user_agent);
}
let protected_headers = if uses_vertex_query_auth {
vec!["content-type"]
+1 -1
View File
@@ -92,7 +92,7 @@ mod upstream_admission;
mod usage;
mod video_tasks;
mod wallet_runtime;
mod xai_profile;
pub use cli_client_profile::ClientProfileSyncGuard;
pub use self::ai_serving::api::{codex_client_originator, codex_client_user_agent};
pub(crate) use self::ai_serving::api::{
+8 -45
View File
@@ -2513,53 +2513,16 @@ async fn run() -> Result<(), Box<dyn std::error::Error>> {
);
}
}
// 两个 CLI 画像的发布检查互不依赖,并发执行以免叠加启动阶段的网络超时。
let (codex_profile, claude_code_profile) = tokio::join!(
state.prewarm_codex_client_profile(),
state.prewarm_claude_code_client_profile(),
);
match codex_profile {
Ok(version) => {
info!(
codex_client_version = %version,
"prewarmed Codex client profile"
);
}
Err(err) => {
warn!(
error = %err,
"failed to refresh Codex client profile; built-in or cached profile remains active"
);
}
}
match claude_code_profile {
Ok(version) => {
info!(
claude_code_client_version = %version,
"prewarmed Claude Code client profile"
);
}
Err(err) => {
warn!(
error = %err,
"failed to refresh Claude Code client profile; built-in or cached profile remains active"
);
}
}
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"
);
for (client, result) in state.prewarm_client_profiles().await {
match result {
Ok(version) => info!(client, version = %version, "prewarmed client profile"),
Err(error) => warn!(client, error = %error,
"client profile refresh failed; built-in or cached profile remains active"),
}
}
// All roles synchronize local snapshots, not just the singleton owner.
// Keep the guard alive until main exits so shutdown cancels the task.
let _client_profile_cache_sync = state.spawn_client_profile_cache_sync();
match prewarm_direct_h2c_sender_cache_from_env_for_startup().await {
Ok(Some(report)) => {
if report.failed_targets > 0 {
+26 -21
View File
@@ -55,7 +55,8 @@ use super::super::{control::GatewayControlDecision, error::GatewayError};
use super::super::{provider_transport, usage};
use crate::cli_client_profile::{
spawn_worker as spawn_cli_client_profile_worker, CLAUDE_CODE_CLI_PROFILE, CODEX_CLI_PROFILE,
spawn_worker as spawn_cli_client_profile_worker, CLAUDE_CODE_CLI_PROFILE, CLI_PROFILES,
CODEX_CLI_PROFILE, XAI_CLI_PROFILE,
};
use crate::maintenance::spawn_account_self_check_worker;
use crate::maintenance::spawn_audit_cleanup_worker;
@@ -78,7 +79,6 @@ 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
@@ -162,7 +162,21 @@ impl AppState {
}
pub async fn prewarm_xai_client_profile(&self) -> Result<String, String> {
crate::xai_profile::prewarm(self.runtime_state()).await
crate::cli_client_profile::prewarm(&XAI_CLI_PROFILE, self.runtime_state()).await
}
pub async fn prewarm_client_profiles(&self) -> Vec<(&'static str, Result<String, String>)> {
futures_util::future::join_all(CLI_PROFILES.iter().map(|spec| async move {
(
spec.client_name(),
crate::cli_client_profile::prewarm(spec, self.runtime_state()).await,
)
}))
.await
}
pub fn spawn_client_profile_cache_sync(&self) -> crate::ClientProfileSyncGuard {
crate::cli_client_profile::spawn_cache_sync(self.clone())
}
pub async fn prewarm_chat_pii_redaction_runtime_config(&self) -> Result<bool, String> {
@@ -2363,24 +2377,15 @@ impl AppState {
crate::task_runtime::TASK_KEY_MODEL_FETCH_WORKER,
spawn_model_fetch_worker(background_state.clone()),
);
supervise_worker(
crate::task_runtime::TASK_KEY_CODEX_CLIENT_PROFILE,
Some(spawn_cli_client_profile_worker(
&CODEX_CLI_PROFILE,
background_state.clone(),
)),
);
supervise_worker(
crate::task_runtime::TASK_KEY_CLAUDE_CODE_CLIENT_PROFILE,
Some(spawn_cli_client_profile_worker(
&CLAUDE_CODE_CLI_PROFILE,
background_state.clone(),
)),
);
supervise_worker(
crate::task_runtime::TASK_KEY_XAI_CLIENT_PROFILE,
Some(spawn_xai_client_profile_worker(background_state.clone())),
);
for spec in CLI_PROFILES {
supervise_worker(
spec.task_key(),
Some(spawn_cli_client_profile_worker(
spec,
background_state.clone(),
)),
);
}
supervise_worker(
crate::task_runtime::TASK_KEY_VIDEO_TASK_POLLER,
spawn_video_task_poller(background_state.clone()),
@@ -28,6 +28,7 @@ pub(crate) const TASK_KEY_CODEX_CLIENT_PROFILE: &str = "maintenance.codex.client
pub(crate) const TASK_KEY_CLAUDE_CODE_CLIENT_PROFILE: &str =
"maintenance.claude_code.client.profile";
pub(crate) const TASK_KEY_XAI_CLIENT_PROFILE: &str = "maintenance.xai.client.profile";
pub(crate) const TASK_KEY_GEMINI_CLI_CLIENT_PROFILE: &str = "maintenance.gemini_cli.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";
@@ -222,6 +223,14 @@ const TASK_DEFINITIONS: &[TaskDefinition] = &[
true,
RETRY_ONCE,
),
TaskDefinition::new(
TASK_KEY_GEMINI_CLI_CLIENT_PROFILE,
TaskKind::Scheduled,
"daily",
true,
true,
RETRY_ONCE,
),
TaskDefinition::new(
TASK_KEY_XAI_CLIENT_PROFILE,
TaskKind::Scheduled,
@@ -978,6 +987,7 @@ mod client_profile_task_tests {
(TASK_KEY_CODEX_CLIENT_PROFILE, "daily"),
(TASK_KEY_CLAUDE_CODE_CLIENT_PROFILE, "daily"),
(TASK_KEY_XAI_CLIENT_PROFILE, "interval"),
(TASK_KEY_GEMINI_CLI_CLIENT_PROFILE, "daily"),
] {
let definitions: Vec<_> = task_definitions()
.iter()
-592
View File
@@ -1,592 +0,0 @@
//! 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");
}
}