chore(release): prepare v0.18.6

This commit is contained in:
晴天
2026-07-14 23:25:42 -07:00
parent 46c9c7743b
commit 90df159d5e
39 changed files with 3349 additions and 909 deletions

2
src-tauri/Cargo.lock generated
View File

@@ -366,7 +366,7 @@ dependencies = [
[[package]]
name = "clawpanel"
version = "0.18.5"
version = "0.18.6"
dependencies = [
"base64 0.22.1",
"chrono",

View File

@@ -1,6 +1,6 @@
[package]
name = "clawpanel"
version = "0.18.5"
version = "0.18.6"
edition = "2021"
description = "ClawPanel - OpenClaw 可视化管理面板"
authors = ["qingchencloud"]

File diff suppressed because it is too large Load Diff

View File

@@ -6,6 +6,7 @@
//! 3. 写入 Hermes Home 下的 config.yaml + .env
use serde_json::Value;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
@@ -1681,12 +1682,12 @@ pub async fn hermes_dashboard_start(app: tauri::AppHandle) -> Result<Value, Stri
// 2. 清掉残留 PID来自上一次 spawn
let _ = kill_dashboard_pid();
// 2b. 上游 wheel 漏装 dashboard_auth + web_dist 时dashboard 进程会直接崩。
// 在 spawn 前先做一次幂等 stub 注入,覆盖既有用户(从早期版本升上来、没走过 install_hermes
// 的新代码路径)也能立即恢复。已存在的真实文件不会被覆盖。
// 2b. 旧版安装缺失 dashboard_auth 时先幂等补齐;最新版源码标签不附带 SPA
// ClawPanel 会另外通过 HERMES_WEB_DIST 提供仅用于 token bootstrap 的最小页面。
inject_hermes_dashboard_compat_stub(&app);
let home = hermes_home();
let dashboard_dist = ensure_hermes_dashboard_fallback_dist(&home)?;
let log_path = home.join("dashboard-run.log");
let log_file =
std::fs::File::create(&log_path).map_err(|e| format!("创建日志文件失败: {e}"))?;
@@ -1696,11 +1697,12 @@ pub async fn hermes_dashboard_start(app: tauri::AppHandle) -> Result<Value, Stri
let enhanced = hermes_enhanced_path();
let mut cmd = std::process::Command::new(hermes_program_for_spawn()?);
cmd.args(["dashboard"])
cmd.args(["dashboard", "--no-open"])
.current_dir(&home)
.stdin(std::process::Stdio::null())
.stdout(log_file)
.stderr(log_err);
cmd.env("HERMES_WEB_DIST", dashboard_dist);
apply_hermes_runtime_env(&mut cmd, &enhanced);
#[cfg(target_os = "windows")]
cmd.creation_flags(CREATE_NO_WINDOW);
@@ -1868,7 +1870,7 @@ pub async fn install_hermes(
let _ = app.emit("hermes-install-progress", 90u32);
// Step 2b: 注入 dashboard 兼容 stub弥补上游 wheel 漏装 dashboard_auth + web_dist
// Step 2b: 为早期缺失 dashboard_auth 的安装保留兼容 stub。
inject_hermes_dashboard_compat_stub(&app);
// Step 3: 验证安装
@@ -1899,7 +1901,40 @@ pub async fn install_hermes(
}
}
/// 确保 uv 二进制可用,不存在则下载
const HERMES_UV_VERSION: &str = "0.11.28";
const HERMES_MIN_UV_VERSION: &str = "0.11.24";
fn parse_version_triplet(text: &str) -> Option<[u64; 3]> {
text.split_whitespace().find_map(|part| {
let normalized = part.trim_matches(|ch: char| !ch.is_ascii_digit() && ch != '.');
let values = normalized
.split('.')
.map(str::parse::<u64>)
.collect::<Result<Vec<_>, _>>()
.ok()?;
(values.len() == 3).then(|| [values[0], values[1], values[2]])
})
}
fn is_hermes_uv_version_supported(version_output: &str) -> bool {
let Some(current) = parse_version_triplet(version_output) else {
return false;
};
let Some(minimum) = parse_version_triplet(HERMES_MIN_UV_VERSION) else {
return false;
};
current >= minimum
}
fn is_uv_wheel_cache_error(text: &str) -> bool {
let lower = text.to_ascii_lowercase();
lower.contains("the wheel is invalid")
|| lower.contains("metadata field name not found")
|| lower.contains("failed to read from the distribution cache")
|| lower.contains("failed to fetch wheel")
}
/// 确保 uv 二进制可用;版本低于缓存冲突修复基线时下载受管版本。
async fn ensure_uv(app: &tauri::AppHandle) -> Result<String, String> {
ensure_hermes_portable_dirs()?;
let uv_path = uv_bin_path();
@@ -1909,8 +1944,16 @@ async fn ensure_uv(app: &tauri::AppHandle) -> Result<String, String> {
if uv_path.exists() {
let path_str = uv_path.to_string_lossy().to_string();
if let Ok(ver) = run_silent(&path_str, &["--version"]) {
let _ = app.emit("hermes-install-log", format!("✓ uv 已就绪: {ver}"));
return Ok(path_str);
if is_hermes_uv_version_supported(&ver) {
let _ = app.emit("hermes-install-log", format!("✓ uv 已就绪: {ver}"));
return Ok(path_str);
}
let _ = app.emit(
"hermes-install-log",
format!(
"⚠ uv 版本过低: {ver}Hermes 安装要求 >= {HERMES_MIN_UV_VERSION},将更新受管 uv"
),
);
}
}
@@ -1923,11 +1966,19 @@ async fn ensure_uv(app: &tauri::AppHandle) -> Result<String, String> {
// 系统 PATH 中有 uv
let enhanced = hermes_enhanced_path();
if let Ok(ver) = run_at_path("uv", &["--version"], &enhanced) {
let _ = app.emit("hermes-install-log", format!("✓ 系统 uv 已就绪: {ver}"));
if let Some(path) = find_executable_path("uv", &enhanced) {
return Ok(path);
if is_hermes_uv_version_supported(&ver) {
let _ = app.emit("hermes-install-log", format!("✓ 系统 uv 已就绪: {ver}"));
if let Some(path) = find_executable_path("uv", &enhanced) {
return Ok(path);
}
return Ok("uv".into());
}
return Ok("uv".into());
let _ = app.emit(
"hermes-install-log",
format!(
"⚠ 系统 {ver} 早于缓存修复版本 {HERMES_MIN_UV_VERSION},将安装受管 uv {HERMES_UV_VERSION}"
),
);
}
}
@@ -1935,7 +1986,7 @@ async fn ensure_uv(app: &tauri::AppHandle) -> Result<String, String> {
let _ = app.emit("hermes-install-log", "📦 下载 uv 包管理器...");
let _ = app.emit("hermes-install-progress", 5u32);
let version = "0.7.12"; // 稳定版本
let version = HERMES_UV_VERSION;
let url = uv_download_url(version);
let _ = app.emit("hermes-install-log", format!("下载: {url}"));
@@ -2051,10 +2102,41 @@ fn extract_uv_tar_gz(data: &[u8], dest: &std::path::Path) -> Result<(), String>
Err("tar.gz 中未找到 uv".into())
}
const HERMES_STABLE_VERSION: &str = "0.18.0";
const HERMES_STABLE_TAG: &str = "v2026.7.1";
const HERMES_STABLE_VERSION: &str = "0.18.2";
const HERMES_STABLE_TAG: &str = "v2026.7.7.2";
const HERMES_GIT_REPO_URL: &str = "https://github.com/NousResearch/hermes-agent.git";
#[cfg(test)]
mod hermes_uv_install_tests {
use super::*;
#[test]
fn rejects_uv_before_archive_collision_fix() {
assert!(!is_hermes_uv_version_supported(
"uv 0.11.14 (3fdfdc7d4 2026-05-12 x86_64-pc-windows-msvc)"
));
assert!(is_hermes_uv_version_supported(
"uv 0.11.24 (5e04460 2026-06-23 x86_64-pc-windows-msvc)"
));
assert!(is_hermes_uv_version_supported(
"uv 0.11.28 (ebf0f43 2026-07-07 x86_64-pc-windows-msvc)"
));
}
#[test]
fn recognizes_invalid_wheel_metadata_as_cache_recoverable() {
assert!(is_uv_wheel_cache_error(
"The wheel is invalid: Metadata field Name not found"
));
assert!(is_uv_wheel_cache_error(
"Failed to read from the distribution cache"
));
assert!(!is_uv_wheel_cache_error(
"Could not resolve host: github.com"
));
}
}
/// Runtime Python deps that `hermes-agent` needs at runtime but are NOT declared as
/// install-time dependencies in its `[project].dependencies` (e.g. lazy-loaded
/// platform adapters). Without these, `hermes gateway run` starts but cannot bring
@@ -2089,8 +2171,10 @@ fn hermes_runtime_extras_log_segment() -> String {
// particular, the missing dashboard_auth subpackage breaks `hermes dashboard`
// completely, taking down every ClawPanel page that talks to port 9119
// (Profile, Kanban, OAuth, Channels, Sessions detail).
// Current stable Hermes v0.18.0 / v2026.7.1 ships these files; this remains a
// no-op compatibility fallback for users upgrading from older broken installs.
// Current stable Hermes v0.18.2 / v2026.7.7.2 ships dashboard_auth, but a source
// install still omits web_dist. ClawPanel-spawned Dashboard processes use the
// managed HERMES_WEB_DIST below; package injection remains as an idempotent
// compatibility fallback for existing and externally started installations.
//
// To stay self-sufficient (per project policy: do not patch upstream), we
// inject a minimal pass-through stub into the installed venv:
@@ -2246,6 +2330,18 @@ const HERMES_DASHBOARD_WEB_DIST_INDEX_HTML: &str = r#"<!doctype html>
</html>
"#;
fn ensure_hermes_dashboard_fallback_dist(home: &Path) -> Result<PathBuf, String> {
let dist = home.join("clawpanel-dashboard-web-dist");
std::fs::create_dir_all(dist.join("assets"))
.map_err(|e| format!("创建 Hermes Dashboard 兼容资源目录失败: {e}"))?;
let index_path = dist.join("index.html");
if !index_path.exists() {
std::fs::write(&index_path, HERMES_DASHBOARD_WEB_DIST_INDEX_HTML)
.map_err(|e| format!("写入 Hermes Dashboard 兼容首页失败: {e}"))?;
}
Ok(dist)
}
/// Resolve `<uv tool dir>/hermes-agent` — the venv root that `uv tool install`
/// creates. Returns `None` if `uv` is unavailable or hermes-agent isn't installed
/// via the uv-tool path (e.g. user is on the legacy `~/.hermes-venv` uv-pip path).
@@ -2632,26 +2728,25 @@ async fn install_via_uv_tool(
let pkg = hermes_package_spec(extras);
let mut cmd = tokio::process::Command::new(uv_path);
cmd.args(["tool", "install", "--force", &pkg, "--python", "3.11"]);
append_hermes_runtime_extras(&mut cmd);
// 配置 PyPI 镜像extras 的依赖仍从 PyPI 下载)
if let Some(mirror) = pypi_mirror_url() {
cmd.args(["--index-url", &mirror]);
}
// 代理
super::apply_proxy_env_tokio(&mut cmd);
let enhanced = hermes_enhanced_path();
apply_hermes_runtime_env_tokio(&mut cmd, &enhanced);
// uv 需要 git 来克隆仓库
cmd.env("GIT_TERMINAL_PROMPT", "0");
// 用户配置了 Git 镜像(如 ghproxy→ 进程级注入 insteadOf 重写
apply_git_mirror_env(&mut cmd);
#[cfg(target_os = "windows")]
cmd.creation_flags(CREATE_NO_WINDOW);
let build_install_command = |no_cache: bool| {
let mut cmd = tokio::process::Command::new(uv_path);
cmd.args(["tool", "install", "--force", &pkg, "--python", "3.11"]);
append_hermes_runtime_extras(&mut cmd);
if no_cache {
cmd.arg("--no-cache");
}
if let Some(mirror) = pypi_mirror_url() {
cmd.args(["--index-url", &mirror]);
}
super::apply_proxy_env_tokio(&mut cmd);
let enhanced = hermes_enhanced_path();
apply_hermes_runtime_env_tokio(&mut cmd, &enhanced);
cmd.env("GIT_TERMINAL_PROMPT", "0");
apply_git_mirror_env(&mut cmd);
#[cfg(target_os = "windows")]
cmd.creation_flags(CREATE_NO_WINDOW);
cmd
};
let _ = app.emit(
"hermes-install-log",
@@ -2662,7 +2757,16 @@ async fn install_via_uv_tool(
);
// 流式执行:安装输出逐行实时显示,不再等进程结束
let (status, stderr_text) = run_install_command_streaming(app, cmd).await?;
let (mut status, mut stderr_text) =
run_install_command_streaming(app, build_install_command(false)).await?;
if !status.success() && is_uv_wheel_cache_error(&stderr_text) {
let _ = app.emit(
"hermes-install-log",
"⚠ 检测到 uv wheel 缓存异常,正在使用隔离缓存自动重试...",
);
(status, stderr_text) =
run_install_command_streaming(app, build_install_command(true)).await?;
}
if status.success() {
let _ = app.emit("hermes-install-log", "✓ uv tool install 完成");
@@ -2747,24 +2851,43 @@ async fn install_via_uv_pip(
let pkg = hermes_package_spec(extras);
let _ = app.emit(
"hermes-install-log",
format!("> uv pip install hermes-agent@{HERMES_STABLE_TAG}"),
format!(
"> uv pip install hermes-agent@{HERMES_STABLE_TAG} {}",
HERMES_RUNTIME_EXTRA_DEPS.join(" ")
),
);
let mut pip_cmd = tokio::process::Command::new(uv_path);
pip_cmd.args(["pip", "install", &pkg]);
pip_cmd.env("GIT_TERMINAL_PROMPT", "0");
pip_cmd.env("VIRTUAL_ENV", &venv_str);
apply_hermes_runtime_env_tokio(&mut pip_cmd, &enhanced);
if let Some(mirror) = pypi_mirror_url() {
pip_cmd.args(["--index-url", &mirror]);
}
apply_git_mirror_env(&mut pip_cmd);
super::apply_proxy_env_tokio(&mut pip_cmd);
#[cfg(target_os = "windows")]
pip_cmd.creation_flags(CREATE_NO_WINDOW);
let build_pip_command = |no_cache: bool| {
let mut pip_cmd = tokio::process::Command::new(uv_path);
pip_cmd.args(["pip", "install", &pkg]);
pip_cmd.args(HERMES_RUNTIME_EXTRA_DEPS);
if no_cache {
pip_cmd.arg("--no-cache");
}
pip_cmd.env("GIT_TERMINAL_PROMPT", "0");
pip_cmd.env("VIRTUAL_ENV", &venv_str);
apply_hermes_runtime_env_tokio(&mut pip_cmd, &enhanced);
if let Some(mirror) = pypi_mirror_url() {
pip_cmd.args(["--index-url", &mirror]);
}
apply_git_mirror_env(&mut pip_cmd);
super::apply_proxy_env_tokio(&mut pip_cmd);
#[cfg(target_os = "windows")]
pip_cmd.creation_flags(CREATE_NO_WINDOW);
pip_cmd
};
// 流式执行pip 下载/构建输出逐行实时显示
let (pip_status, pip_stderr) = run_install_command_streaming(app, pip_cmd).await?;
let (mut pip_status, mut pip_stderr) =
run_install_command_streaming(app, build_pip_command(false)).await?;
if !pip_status.success() && is_uv_wheel_cache_error(&pip_stderr) {
let _ = app.emit(
"hermes-install-log",
"⚠ 检测到 uv wheel 缓存异常,正在使用隔离缓存自动重试...",
);
(pip_status, pip_stderr) =
run_install_command_streaming(app, build_pip_command(true)).await?;
}
if !pip_status.success() {
let cleaned = sanitize_hermes_install_output(pip_stderr.trim());
@@ -3104,31 +3227,104 @@ fn merge_env_file(existing: &str, managed_keys: &[&str], new_pairs: &[(String, S
content
}
fn hermes_stable_backup_path(file: &Path) -> PathBuf {
let mut backup = file.as_os_str().to_os_string();
backup.push(".bak");
PathBuf::from(backup)
}
fn replace_hermes_files_transaction(entries: &[(PathBuf, String, bool)]) -> Result<(), String> {
let suffix = format!(
"{}-{}",
std::process::id(),
chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default()
);
let mut staged: Vec<(PathBuf, PathBuf, PathBuf, bool, bool)> = Vec::new();
let mut staged: Vec<(PathBuf, PathBuf, PathBuf, bool, bool, String)> = Vec::new();
for (file, content, private) in entries {
#[cfg(not(unix))]
let _ = private;
if let Some(parent) = file.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("创建 Hermes 配置目录失败: {e}"))?;
let staging_result = (|| {
for (file, content, private) in entries {
#[cfg(not(unix))]
let _ = private;
if let Some(parent) = file.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("创建 Hermes 配置目录失败: {e}"))?;
}
let temp = file.with_extension(format!("tmp-{suffix}"));
let backup = file.with_extension(format!("bak-sync-{suffix}"));
let mut staged_file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&temp)
.map_err(|e| format!("创建 Hermes 临时配置失败: {e}"))?;
if let Err(error) = staged_file
.write_all(content.as_bytes())
.and_then(|_| staged_file.sync_all())
{
drop(staged_file);
let _ = std::fs::remove_file(&temp);
return Err(format!("写入 Hermes 临时配置失败: {error}"));
}
drop(staged_file);
#[cfg(unix)]
if *private {
use std::os::unix::fs::PermissionsExt;
if let Err(error) =
std::fs::set_permissions(&temp, std::fs::Permissions::from_mode(0o600))
{
let _ = std::fs::remove_file(&temp);
return Err(format!("设置临时配置权限失败: {error}"));
}
}
staged.push((
file.clone(),
temp,
backup,
file.exists(),
false,
content.clone(),
));
}
let temp = file.with_extension(format!("tmp-{suffix}"));
let backup = file.with_extension(format!("bak-sync-{suffix}"));
std::fs::write(&temp, content).map_err(|e| format!("写入临时配置失败: {e}"))?;
#[cfg(unix)]
if *private {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&temp, std::fs::Permissions::from_mode(0o600))
.map_err(|e| format!("设置临时配置权限失败: {e}"))?;
Ok(())
})();
if let Err(error) = staging_result {
for entry in &staged {
let _ = std::fs::remove_file(&entry.1);
let _ = std::fs::remove_file(&entry.2);
}
staged.push((file.clone(), temp, backup, file.exists(), false));
return Err(error);
}
let stable_backup_result = (|| {
for (file, _, private) in entries {
#[cfg(not(unix))]
let _ = private;
if !file.exists() {
continue;
}
let stable_backup = hermes_stable_backup_path(file);
std::fs::copy(file, &stable_backup)
.map_err(|e| format!("创建 Hermes 稳定备份失败: {e}"))?;
std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&stable_backup)
.and_then(|backup| backup.sync_all())
.map_err(|e| format!("同步 Hermes 稳定备份失败: {e}"))?;
#[cfg(unix)]
if *private {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&stable_backup, std::fs::Permissions::from_mode(0o600))
.map_err(|e| format!("设置 Hermes 稳定备份权限失败: {e}"))?;
}
}
Ok(())
})();
if let Err(error) = stable_backup_result {
for entry in &staged {
let _ = std::fs::remove_file(&entry.1);
let _ = std::fs::remove_file(&entry.2);
}
return Err(error);
}
let result = (|| {
@@ -3141,6 +3337,13 @@ fn replace_hermes_files_transaction(entries: &[(PathBuf, String, bool)]) -> Resu
.map_err(|e| format!("提交 Hermes 配置失败: {e}"))?;
entry.4 = true;
}
for entry in &staged {
let readback = std::fs::read_to_string(&entry.0)
.map_err(|e| format!("回读 Hermes 配置失败: {e}"))?;
if readback != entry.5 {
return Err("Hermes 配置写入后回读不一致".into());
}
}
Ok(())
})();
@@ -3231,7 +3434,7 @@ fn sync_hermes_provider_files_at(
}
replace_hermes_files_transaction(&entries)?;
Ok(serde_json::json!({ "providerId": provider, "envKey": env_key }))
Ok(serde_json::json!({ "providerId": provider, "envKey": env_key, "verified": true }))
}
#[tauri::command]
@@ -3242,6 +3445,7 @@ pub fn hermes_sync_provider(
model: Option<String>,
set_default: bool,
) -> Result<Value, String> {
let api_key = super::config::resolve_model_api_key(&api_key)?;
sync_hermes_provider_files_at(
&hermes_home(),
&provider,
@@ -26197,7 +26401,7 @@ mod hermes_provider_sync_tests {
)
.unwrap();
sync_hermes_provider_files_at(
let result = sync_hermes_provider_files_at(
&home,
"custom",
"sk-new",
@@ -26206,6 +26410,7 @@ mod hermes_provider_sync_tests {
true,
)
.unwrap();
assert_eq!(result["verified"], serde_json::json!(true));
let env = std::fs::read_to_string(home.join(".env")).unwrap();
assert!(env.contains("ANTHROPIC_API_KEY=keep-me"));
@@ -26218,6 +26423,14 @@ mod hermes_provider_sync_tests {
assert!(config.contains("default: gpt-test"));
assert!(config.contains("provider: custom"));
assert!(config.contains("level: INFO"));
assert_eq!(
std::fs::read_to_string(home.join(".env.bak")).unwrap(),
"ANTHROPIC_API_KEY=keep-me\nOPENAI_API_KEY=old\nCUSTOM_FLAG=keep\n"
);
assert_eq!(
std::fs::read_to_string(home.join("config.yaml.bak")).unwrap(),
"model:\n default: old-model\n provider: anthropic\nlogging:\n level: INFO\n"
);
let _ = std::fs::remove_dir_all(&home);
}

View File

@@ -1,13 +1,17 @@
use base64::{engine::general_purpose, Engine as _};
use futures_util::StreamExt;
use serde_json::{json, Map, Value};
use std::collections::HashSet;
use std::path::{Component, Path, PathBuf};
use std::sync::Mutex;
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use tokio::io::AsyncWriteExt;
use tokio::sync::Semaphore;
/// media-jobs.json 的读改写锁:并发轮询/写入时防止互相覆盖
static MEDIA_JOBS_LOCK: Mutex<()> = Mutex::new(());
static MEDIA_QUEUE_INFLIGHT: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
static MEDIA_QUEUE_SEMAPHORE: OnceLock<Arc<Semaphore>> = OnceLock::new();
fn lock_media_jobs() -> std::sync::MutexGuard<'static, ()> {
MEDIA_JOBS_LOCK
@@ -25,6 +29,8 @@ const MEDIA_CONFIG_FILE: &str = "media-config.json";
const MEDIA_JOBS_FILE: &str = "media-jobs.json";
const MAX_IMAGE_COUNT: u64 = 4;
const MAX_ASSET_BYTES: u64 = 512 * 1024 * 1024;
const MEDIA_QUEUE_CONCURRENCY: usize = 2;
const MEDIA_QUEUE_POLL_INTERVAL: Duration = Duration::from_secs(5);
/// 内嵌预览base64 IPC上限超过此大小引导用户打开文件夹本地查看
const MAX_INLINE_PREVIEW_BYTES: u64 = 64 * 1024 * 1024;
@@ -366,16 +372,95 @@ fn read_json_or_default(path: &Path, default: Value) -> Value {
}
pub(super) fn write_json_atomic(path: &Path, value: &Value) -> Result<(), String> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| format!("创建目录失败: {e}"))?;
}
let parent = path
.parent()
.ok_or_else(|| "JSON 路径缺少父目录".to_string())?;
std::fs::create_dir_all(parent).map_err(|e| format!("创建目录失败: {e}"))?;
let content = serde_json::to_string_pretty(value).map_err(|e| format!("序列化失败: {e}"))?;
let tmp = path.with_extension("tmp");
std::fs::write(&tmp, content).map_err(|e| format!("写入临时文件失败: {e}"))?;
if path.exists() {
let _ = std::fs::remove_file(path);
let parsed: Value =
serde_json::from_str(&content).map_err(|e| format!("候选 JSON 校验失败: {e}"))?;
if &parsed != value {
return Err("候选 JSON 序列化后内容不一致".into());
}
std::fs::rename(&tmp, path).map_err(|e| format!("替换文件失败: {e}"))
let suffix = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("data.json");
let tmp = parent.join(format!(".{name}.{}.{suffix}.tmp", std::process::id()));
let rollback = parent.join(format!(".{name}.{}.{suffix}.old", std::process::id()));
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&tmp)
.map_err(|e| format!("创建临时文件失败: {e}"))?;
if let Err(error) =
std::io::Write::write_all(&mut file, content.as_bytes()).and_then(|_| file.sync_all())
{
let _ = std::fs::remove_file(&tmp);
return Err(format!("写入临时文件失败: {error}"));
}
drop(file);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
if let Err(error) = std::fs::set_permissions(&tmp, std::fs::Permissions::from_mode(0o600)) {
let _ = std::fs::remove_file(&tmp);
return Err(format!("设置临时文件权限失败: {error}"));
}
}
let had_existing = path.exists();
#[cfg(not(target_os = "windows"))]
let replace_result = {
if had_existing {
std::fs::copy(path, &rollback).map(|_| ())
} else {
Ok(())
}
.and_then(|_| std::fs::rename(&tmp, path))
};
#[cfg(target_os = "windows")]
let replace_result = if had_existing {
std::fs::rename(path, &rollback).and_then(|_| {
std::fs::rename(&tmp, path).inspect_err(|_| {
let _ = std::fs::rename(&rollback, path);
})
})
} else {
std::fs::rename(&tmp, path)
};
if let Err(error) = replace_result {
let _ = std::fs::remove_file(&tmp);
let _ = std::fs::remove_file(&rollback);
return Err(format!("替换 JSON 文件失败,原文件已保留: {error}"));
}
let readback = std::fs::read_to_string(path)
.ok()
.and_then(|raw| serde_json::from_str::<Value>(&raw).ok());
if readback.as_ref() != Some(value) {
if had_existing && rollback.exists() {
let _ = std::fs::copy(&rollback, path);
} else {
let _ = std::fs::remove_file(path);
}
let _ = std::fs::remove_file(&rollback);
return Err("JSON 写入后回读不一致,已恢复原文件".into());
}
let _ = std::fs::remove_file(&rollback);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600));
}
Ok(())
}
fn read_media_config_private() -> Value {
@@ -465,6 +550,57 @@ fn new_job_id(prefix: &str) -> String {
)
}
fn build_queued_media_job(
job_id: &str,
kind: &str,
provider: &MediaProviderConfig,
request: &Value,
) -> Value {
let now = now_iso();
let model = if kind == "video" {
&provider.video_model
} else {
&provider.image_model
};
json!({
"id": job_id,
"type": kind,
"provider": provider.provider,
"providerTaskId": null,
"status": "queued",
"providerStatus": "local-queued",
"prompt": str_field(request, "prompt"),
"model": model,
"createdAt": now,
"updatedAt": now,
"request": truncate_large_strings(request),
"assets": [],
"error": null,
"rawProviderResponse": null
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum MediaRecoveryAction {
Schedule,
FailUnknown,
Ignore,
}
fn media_recovery_action(job: &Value) -> MediaRecoveryAction {
match str_field(job, "status") {
"queued" => MediaRecoveryAction::Schedule,
"running"
if str_field(job, "type") == "video"
&& !str_field(job, "providerTaskId").is_empty() =>
{
MediaRecoveryAction::Schedule
}
"running" => MediaRecoveryAction::FailUnknown,
_ => MediaRecoveryAction::Ignore,
}
}
fn validate_provider_id(provider: &str) -> Result<(), String> {
if is_supported_media_provider(provider) {
Ok(())
@@ -1396,32 +1532,6 @@ fn truncate_large_strings(value: &Value) -> Value {
}
}
fn media_error_job(
id: &str,
kind: &str,
provider: &str,
prompt: &str,
model: &str,
error: String,
raw_response: Option<Value>,
) -> Value {
let now = now_iso();
json!({
"id": id,
"type": kind,
"provider": provider,
"providerTaskId": null,
"status": "failed",
"prompt": prompt,
"model": model,
"createdAt": now,
"updatedAt": now,
"assets": [],
"error": error,
"rawProviderResponse": raw_response.map(|v| truncate_large_strings(&v)).unwrap_or(Value::Null)
})
}
#[tauri::command]
pub fn read_media_config() -> Result<Value, String> {
ensure_media_root()?;
@@ -1524,15 +1634,11 @@ pub async fn fetch_media_models(provider: String) -> Result<Value, String> {
}))
}
#[tauri::command]
pub async fn generate_image(request: Value) -> Result<Value, String> {
validate_image_request(&request)?;
let provider_id = str_field(&request, "provider").to_string();
let prompt = str_field(&request, "prompt").to_string();
let provider = load_provider_config(&provider_id, "image", str_field(&request, "model"))?;
let job_id = new_job_id("img");
async fn execute_queued_image_job(job_id: &str, request: &Value) -> Result<Value, String> {
let provider_id = str_field(request, "provider").to_string();
let provider = load_provider_config(&provider_id, "image", str_field(request, "model"))?;
let client = media_http_client(provider.timeout_seconds)?;
let payload = build_image_generation_payload(&provider, &request);
let payload = build_image_generation_payload(&provider, request);
let endpoint = build_api_url(&provider.base_url, "/images/generations");
let resp = client
.post(endpoint)
@@ -1549,63 +1655,28 @@ pub async fn generate_image(request: Value) -> Result<Value, String> {
let parsed = serde_json::from_str::<Value>(&text).unwrap_or_else(|_| json!({ "raw": text }));
if !status.is_success() {
let error = sanitize_provider_error(&format!("HTTP {status}: {parsed}"), &provider.api_key);
let job = media_error_job(
&job_id,
"image",
&provider.provider,
&prompt,
&provider.image_model,
error.clone(),
Some(parsed),
);
let _ = upsert_media_job(job);
return Err(error);
}
let assets = match save_image_outputs(&client, &provider, &job_id, &parsed).await {
Ok(assets) => assets,
Err(error) => {
let job = media_error_job(
&job_id,
"image",
&provider.provider,
&prompt,
&provider.image_model,
error.clone(),
Some(parsed),
);
let _ = upsert_media_job(job);
return Err(error);
let assets = save_image_outputs(&client, &provider, job_id, &parsed).await?;
update_media_job(job_id, |entry| {
entry["request"] = truncate_large_strings(&payload);
entry["assets"] = Value::Array(assets.clone());
entry["rawProviderResponse"] = truncate_large_strings(&parsed);
entry["error"] = Value::Null;
if str_field(entry, "status") != "canceled" {
entry["status"] = Value::String("succeeded".into());
entry["providerStatus"] = Value::String("completed".into());
}
};
let now = now_iso();
let job = json!({
"id": job_id,
"type": "image",
"provider": provider.provider,
"providerTaskId": null,
"status": "succeeded",
"prompt": prompt,
"model": provider.image_model,
"createdAt": now,
"updatedAt": now,
"request": truncate_large_strings(&payload),
"assets": assets,
"error": null,
"rawProviderResponse": truncate_large_strings(&parsed)
});
upsert_media_job(job)
})
}
#[tauri::command]
pub async fn create_video_task(request: Value) -> Result<Value, String> {
validate_video_request(&request)?;
let provider_id = str_field(&request, "provider").to_string();
let prompt = str_field(&request, "prompt").to_string();
let provider = load_provider_config(&provider_id, "video", str_field(&request, "model"))?;
let job_id = new_job_id("vid");
async fn submit_queued_video_job(job_id: &str, request: &Value) -> Result<Value, String> {
let provider_id = str_field(request, "provider").to_string();
let prompt = str_field(request, "prompt").to_string();
let provider = load_provider_config(&provider_id, "video", str_field(request, "model"))?;
let client = media_http_client(provider.timeout_seconds)?;
let (payload, resp) = if is_openai_compatible_provider(&provider.provider) {
let payload = build_openai_video_payload(&provider, &request);
let payload = build_openai_video_payload(&provider, request);
let mut form = reqwest::multipart::Form::new()
.text("model", str_field(&payload, "model").to_string())
.text("prompt", str_field(&payload, "prompt").to_string())
@@ -1647,9 +1718,9 @@ pub async fn create_video_task(request: Value) -> Result<Value, String> {
let mut payload = json!({
"model": provider.video_model.clone(),
"content": content,
"ratio": str_field(&request, "ratio"),
"resolution": str_field(&request, "resolution"),
"duration": u64_field(&request, "duration", 5)
"ratio": str_field(request, "ratio"),
"resolution": str_field(request, "resolution"),
"duration": u64_field(request, "duration", 5)
});
for key in ["ratio", "resolution"] {
if payload[key].as_str().unwrap_or("").is_empty() {
@@ -1674,58 +1745,28 @@ pub async fn create_video_task(request: Value) -> Result<Value, String> {
let parsed = serde_json::from_str::<Value>(&text).unwrap_or_else(|_| json!({ "raw": text }));
if !status.is_success() {
let error = sanitize_provider_error(&format!("HTTP {status}: {parsed}"), &provider.api_key);
let job = media_error_job(
&job_id,
"video",
&provider.provider,
&prompt,
&provider.video_model,
error.clone(),
Some(parsed),
);
let _ = upsert_media_job(job);
return Err(error);
}
let provider_task_id = extract_provider_task_id(&parsed);
if provider_task_id.is_empty() {
// 任务可能已在服务商侧创建并计费:即使无法识别任务 ID也要落盘保留原始响应供恢复
let error = "服务商响应中没有视频任务 ID".to_string();
let _ = upsert_media_job(media_error_job(
&job_id,
"video",
&provider.provider,
&prompt,
&provider.video_model,
error.clone(),
Some(parsed),
));
return Err(error);
return Err("服务商响应中没有视频任务 ID".into());
}
let provider_status = extract_provider_status(&parsed);
let job_status = provider_status_to_job_status(&provider_status);
let now = now_iso();
let job = json!({
"id": job_id,
"type": "video",
"provider": provider.provider,
"providerTaskId": provider_task_id,
"status": job_status,
"providerStatus": provider_status,
"prompt": prompt,
"model": provider.video_model,
"createdAt": now,
"updatedAt": now,
"request": truncate_large_strings(&payload),
"assets": [],
"error": null,
"rawProviderResponse": truncate_large_strings(&parsed)
});
upsert_media_job(job)
update_media_job(job_id, |entry| {
entry["providerTaskId"] = Value::String(provider_task_id.clone());
entry["providerStatus"] = Value::String(provider_status.clone());
if str_field(entry, "status") != "canceled" {
entry["status"] = Value::String(job_status.into());
}
entry["request"] = truncate_large_strings(&payload);
entry["rawProviderResponse"] = truncate_large_strings(&parsed);
entry["error"] = Value::Null;
})
}
#[tauri::command]
pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
let job = media_job_by_id(&job_id).ok_or_else(|| format!("媒体任务不存在: {job_id}"))?;
async fn poll_video_task_internal(job_id: &str) -> Result<Value, String> {
let job = media_job_by_id(job_id).ok_or_else(|| format!("媒体任务不存在: {job_id}"))?;
if str_field(&job, "type") != "video" {
return Err("只支持轮询视频任务".into());
}
@@ -1760,7 +1801,7 @@ pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
let error = sanitize_provider_error(&format!("HTTP {status}: {parsed}"), &provider.api_key);
// 429 / 5xx 属于瞬时错误:保留原状态允许继续轮询,避免把服务商侧仍在运行的付费任务永久标记为失败
let transient = status.as_u16() == 429 || status.is_server_error();
return update_media_job(&job_id, |entry| {
return update_media_job(job_id, |entry| {
if !transient {
entry["status"] = Value::String("failed".into());
}
@@ -1776,7 +1817,7 @@ pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
let urls = collect_video_urls(&parsed);
let mut saved = Vec::new();
for (idx, url) in urls.iter().enumerate() {
match download_asset_to_media_root(&client, &provider, url, "video", &job_id, idx).await
match download_asset_to_media_root(&client, &provider, url, "video", job_id, idx).await
{
Ok(asset) => saved.push(asset),
Err(e) => download_error = Some(e),
@@ -1787,7 +1828,7 @@ pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
&client,
&provider,
&provider_task_id,
&job_id,
job_id,
saved.len(),
)
.await
@@ -1798,7 +1839,7 @@ pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
}
assets = Value::Array(saved);
}
update_media_job(&job_id, |entry| {
update_media_job(job_id, |entry| {
entry["status"] = Value::String(job_status.into());
entry["providerStatus"] = Value::String(provider_status.clone());
entry["assets"] = assets.clone();
@@ -1811,6 +1852,196 @@ pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
})
}
fn media_queue_inflight() -> &'static Mutex<HashSet<String>> {
MEDIA_QUEUE_INFLIGHT.get_or_init(|| Mutex::new(HashSet::new()))
}
fn media_queue_semaphore() -> Arc<Semaphore> {
MEDIA_QUEUE_SEMAPHORE
.get_or_init(|| Arc::new(Semaphore::new(MEDIA_QUEUE_CONCURRENCY)))
.clone()
}
fn media_job_is_inflight(job_id: &str) -> bool {
media_queue_inflight()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(job_id)
}
fn mark_media_job_failed(job_id: &str, error: &str) {
let _ = update_media_job(job_id, |entry| {
if matches!(str_field(entry, "status"), "canceled" | "succeeded") {
return;
}
entry["status"] = Value::String("failed".into());
entry["providerStatus"] = Value::String("local-failed".into());
entry["error"] = Value::String(error.to_string());
});
}
async fn process_media_queue_job(job_id: &str) -> Result<(), String> {
let mut job = media_job_by_id(job_id).ok_or_else(|| format!("媒体任务不存在: {job_id}"))?;
if matches!(
str_field(&job, "status"),
"succeeded" | "failed" | "canceled"
) {
return Ok(());
}
if str_field(&job, "status") == "queued" {
job = update_media_job(job_id, |entry| {
if str_field(entry, "status") != "queued" {
return;
}
entry["status"] = Value::String("running".into());
entry["providerStatus"] = Value::String("local-starting".into());
entry["startedAt"] = Value::String(now_iso());
entry["error"] = Value::Null;
})?;
if str_field(&job, "status") != "running" {
return Ok(());
}
let request = job.get("request").cloned().unwrap_or_else(|| json!({}));
if str_field(&job, "type") == "image" {
execute_queued_image_job(job_id, &request).await?;
return Ok(());
}
submit_queued_video_job(job_id, &request).await?;
}
loop {
job = match media_job_by_id(job_id) {
Some(job) => job,
None => return Ok(()),
};
if matches!(str_field(&job, "status"), "canceled" | "failed") {
return Ok(());
}
if str_field(&job, "status") == "succeeded"
&& job
.get("assets")
.and_then(Value::as_array)
.is_some_and(|assets| !assets.is_empty())
{
return Ok(());
}
if str_field(&job, "providerTaskId").is_empty() {
return Err("视频任务缺少服务商任务 ID".into());
}
match poll_video_task_internal(job_id).await {
Ok(updated)
if matches!(
str_field(&updated, "status"),
"succeeded" | "failed" | "canceled"
) =>
{
return Ok(())
}
Ok(_) => {}
Err(error) => {
let _ = update_media_job(job_id, |entry| {
if str_field(entry, "status") == "running" {
entry["error"] = Value::String(error.clone());
}
});
}
}
tokio::time::sleep(MEDIA_QUEUE_POLL_INTERVAL).await;
}
}
fn schedule_media_job(job_id: String) -> bool {
{
let mut inflight = media_queue_inflight()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if !inflight.insert(job_id.clone()) {
return false;
}
}
tauri::async_runtime::spawn(async move {
let permit = match media_queue_semaphore().acquire_owned().await {
Ok(permit) => permit,
Err(error) => {
mark_media_job_failed(&job_id, &format!("后台队列不可用: {error}"));
media_queue_inflight()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&job_id);
return;
}
};
if let Err(error) = process_media_queue_job(&job_id).await {
mark_media_job_failed(&job_id, &error);
}
drop(permit);
media_queue_inflight()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&job_id);
});
true
}
fn recover_media_queue() {
let jobs = read_media_jobs_private()
.get("jobs")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
for job in jobs {
let job_id = str_field(&job, "id").to_string();
if job_id.is_empty() || media_job_is_inflight(&job_id) {
continue;
}
match media_recovery_action(&job) {
MediaRecoveryAction::Schedule => {
schedule_media_job(job_id);
}
MediaRecoveryAction::FailUnknown => mark_media_job_failed(
&job_id,
"应用重启时任务仍在提交中,服务商状态未知,请确认账单后重试",
),
MediaRecoveryAction::Ignore => {}
}
}
}
#[tauri::command]
pub async fn generate_image(request: Value) -> Result<Value, String> {
validate_image_request(&request)?;
let provider_id = str_field(&request, "provider");
let provider = load_provider_config(provider_id, "image", str_field(&request, "model"))?;
let job_id = new_job_id("img");
let job = upsert_media_job(build_queued_media_job(
&job_id, "image", &provider, &request,
))?;
schedule_media_job(job_id);
Ok(job)
}
#[tauri::command]
pub async fn create_video_task(request: Value) -> Result<Value, String> {
validate_video_request(&request)?;
let provider_id = str_field(&request, "provider");
let provider = load_provider_config(provider_id, "video", str_field(&request, "model"))?;
let job_id = new_job_id("vid");
let job = upsert_media_job(build_queued_media_job(
&job_id, "video", &provider, &request,
))?;
schedule_media_job(job_id);
Ok(job)
}
#[tauri::command]
pub async fn poll_video_task(job_id: String) -> Result<Value, String> {
let job = poll_video_task_internal(&job_id).await?;
if str_field(&job, "status") == "running" {
schedule_media_job(job_id);
}
Ok(job)
}
#[tauri::command]
pub async fn cancel_media_job(job_id: String) -> Result<Value, String> {
update_media_job(&job_id, |entry| {
@@ -1821,6 +2052,7 @@ pub async fn cancel_media_job(job_id: String) -> Result<Value, String> {
#[tauri::command]
pub fn list_media_jobs(filter: Option<Value>) -> Result<Value, String> {
recover_media_queue();
let filter = filter.unwrap_or_else(|| json!({}));
let jobs_doc = read_media_jobs_private();
Ok(media_jobs_response_from_doc(&jobs_doc, Some(&filter)))
@@ -2114,6 +2346,55 @@ mod tests {
assert_eq!(provider_status_to_job_status("cancelled"), "canceled");
}
#[test]
fn queued_media_job_keeps_request_without_provider_credentials() {
let provider = MediaProviderConfig {
provider: "openai".into(),
base_url: "https://api.openai.com/v1".into(),
api_key: "sk-secret".into(),
image_model: "gpt-image-1".into(),
video_model: "sora-2".into(),
timeout_seconds: 600,
};
let request = json!({
"provider": "openai",
"prompt": "batch image",
"model": "gpt-image-1",
"count": 2
});
let job = build_queued_media_job("img-test", "image", &provider, &request);
assert_eq!(job["status"], "queued");
assert_eq!(job["request"]["prompt"], "batch image");
assert_eq!(job["request"]["count"], 2);
assert!(job.to_string().find("sk-secret").is_none());
}
#[test]
fn recovery_only_resumes_safe_persisted_jobs() {
assert_eq!(
media_recovery_action(&json!({ "status": "queued", "type": "image" })),
MediaRecoveryAction::Schedule
);
assert_eq!(
media_recovery_action(&json!({
"status": "running",
"type": "video",
"providerTaskId": "provider-task-1"
})),
MediaRecoveryAction::Schedule
);
assert_eq!(
media_recovery_action(&json!({ "status": "running", "type": "image" })),
MediaRecoveryAction::FailUnknown
);
assert_eq!(
media_recovery_action(&json!({ "status": "succeeded", "type": "image" })),
MediaRecoveryAction::Ignore
);
}
#[test]
fn same_host_asset_download_retries_without_auth_on_auth_rejection() {
assert!(should_retry_asset_without_auth(

View File

@@ -13,7 +13,7 @@
//! 存储位置openclaw_dir/clawpanel/model-channels.json —— 跟随 OpenClaw
//! 数据目录,便携迁移整体复制后自动生效(与媒体数据同一决策)。
use serde_json::{json, Map, Value};
use serde_json::{json, Value};
use std::path::PathBuf;
const CHANNELS_FILE: &str = "model-channels.json";
@@ -25,7 +25,7 @@ fn channels_path() -> PathBuf {
}
fn default_channels_doc() -> Value {
json!({ "version": 1, "channels": [], "syncState": {} })
json!({ "version": 2, "channels": [], "syncState": {} })
}
fn str_of(value: &Value, key: &str) -> String {
@@ -41,7 +41,18 @@ fn is_keep_sentinel(key: &str) -> bool {
key.is_empty() || key == "__KEEP__" || key == "••••••••" || key == "********"
}
/// 归一化单个模型条目:接受字符串或 { id, name? } 对象id 为空则丢弃
fn normalize_api_type(raw: &str) -> String {
match raw.trim().to_ascii_lowercase().as_str() {
"" => "openai-completions".into(),
"openai-codex-responses" => "openai-chatgpt-responses".into(),
"google-gemini" | "gemini" | "google" => "google-generative-ai".into(),
"anthropic" => "anthropic-messages".into(),
"openai" | "openai-chat" => "openai-completions".into(),
other => other.to_string(),
}
}
/// 归一化单个模型条目,并保留 OpenClaw 的能力、成本和兼容性元数据。
fn normalize_model_entry(entry: &Value) -> Option<Value> {
if let Some(id) = entry.as_str() {
let id = id.trim();
@@ -54,18 +65,44 @@ fn normalize_model_entry(entry: &Value) -> Option<Value> {
if id.is_empty() {
return None;
}
let mut out = Map::new();
let mut out = entry.as_object()?.clone();
out.insert("id".into(), Value::String(id));
let name = str_of(entry, "name");
if !name.is_empty() {
out.insert("name".into(), Value::String(name));
} else {
out.remove("name");
}
if let Some(ctx) = entry.get("contextWindow").and_then(Value::as_u64) {
out.insert("contextWindow".into(), Value::Number(ctx.into()));
for key in ["contextWindow", "contextTokens", "maxTokens"] {
match entry
.get(key)
.and_then(Value::as_u64)
.filter(|value| *value > 0)
{
Some(value) => {
out.insert(key.into(), Value::Number(value.into()));
}
None => {
out.remove(key);
}
}
}
if let Some(api) = entry.get("api").and_then(Value::as_str) {
out.insert("api".into(), Value::String(normalize_api_type(api)));
}
Some(Value::Object(out))
}
fn normalize_api_key_ref(value: Option<&Value>) -> Option<Value> {
let value = value?.as_object()?;
if str_of(&Value::Object(value.clone()), "source").is_empty()
|| str_of(&Value::Object(value.clone()), "id").is_empty()
{
return None;
}
Some(Value::Object(value.clone()))
}
/// 归一化单个渠道current 为同 id 的旧渠道(用于保留旧 Key
/// 返回 None 表示条目非法(缺 id/名称),直接丢弃。
fn normalize_channel(entry: &Value, current: Option<&Value>) -> Option<Value> {
@@ -81,11 +118,37 @@ fn normalize_channel(entry: &Value, current: Option<&Value>) -> Option<Value> {
}
let incoming_key = str_of(entry, "apiKey");
let api_key = if is_keep_sentinel(&incoming_key) {
current.map(|c| str_of(c, "apiKey")).unwrap_or_default()
let keep_key = is_keep_sentinel(&incoming_key);
let current_key = current.map(|c| str_of(c, "apiKey")).unwrap_or_default();
let api_key = if keep_key {
current_key.clone()
} else {
incoming_key
};
let incoming_ref = normalize_api_key_ref(entry.get("apiKeyRef"));
let current_ref = normalize_api_key_ref(current.and_then(|value| value.get("apiKeyRef")));
let api_key_ref = if keep_key {
incoming_ref.or_else(|| current_ref.clone())
} else {
None
};
let stored_version = entry
.get("credentialVersion")
.and_then(Value::as_u64)
.unwrap_or(0);
let current_version = current
.and_then(|value| value.get("credentialVersion"))
.and_then(Value::as_u64)
.unwrap_or_else(|| u64::from(!current_key.is_empty()));
let credential_version = if current.is_some() {
if (!keep_key && api_key != current_key) || api_key_ref != current_ref {
current_version.saturating_add(1)
} else {
current_version
}
} else {
stored_version.max(u64::from(!api_key.is_empty() || api_key_ref.is_some()))
};
let mut models: Vec<Value> = Vec::new();
let mut seen = std::collections::HashSet::new();
@@ -109,27 +172,34 @@ fn normalize_channel(entry: &Value, current: Option<&Value>) -> Option<Value> {
models.first().map(|m| str_of(m, "id")).unwrap_or_default()
};
let api_type = {
let t = str_of(entry, "apiType");
if t.is_empty() {
"openai-completions".to_string()
} else {
t
}
};
let api_type = normalize_api_type(&str_of(entry, "apiType"));
let mut provider_config = entry
.get("providerConfig")
.and_then(Value::as_object)
.cloned()
.unwrap_or_default();
for managed in ["baseUrl", "api", "apiKey", "models"] {
provider_config.remove(managed);
}
Some(json!({
let mut normalized = json!({
"id": id,
"name": name,
"presetKey": str_of(entry, "presetKey"),
"baseUrl": base_url,
"apiType": api_type,
"apiKey": api_key,
"credentialVersion": credential_version,
"providerConfig": provider_config,
"models": models,
"defaultModel": default_model,
"enabled": entry.get("enabled").and_then(Value::as_bool).unwrap_or(true),
"updatedAt": str_of(entry, "updatedAt"),
}))
});
if let Some(api_key_ref) = api_key_ref {
normalized["apiKeyRef"] = api_key_ref;
}
Some(normalized)
}
/// 归一化整个文档current 提供旧文档以支持保留旧 Key
@@ -161,7 +231,7 @@ fn normalize_channels_doc(config: &Value, current: Option<&Value>) -> Value {
.cloned()
.unwrap_or_else(|| json!({}));
json!({ "version": 1, "channels": channels, "syncState": sync_state })
json!({ "version": 2, "channels": channels, "syncState": sync_state })
}
fn read_channels_private() -> Value {
@@ -182,14 +252,19 @@ fn sanitize_doc_for_read(doc: &Value) -> Value {
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let has_api_key_ref = normalize_api_key_ref(obj.get("apiKeyRef")).is_some();
obj.insert("apiKey".into(), Value::String(String::new()));
obj.insert(
"apiKeySaved".into(),
Value::Bool(!api_key.trim().is_empty()),
Value::Bool(!api_key.trim().is_empty() || has_api_key_ref),
);
obj.insert(
"apiKeyMask".into(),
Value::String(super::media::api_key_mask(&api_key)),
Value::String(if has_api_key_ref {
"SecretRef".into()
} else {
super::media::api_key_mask(&api_key)
}),
);
}
}
@@ -206,7 +281,24 @@ pub fn read_model_channels() -> Result<Value, String> {
pub fn write_model_channels(config: Value) -> Result<Value, String> {
let current = read_channels_private();
let normalized = normalize_channels_doc(&config, Some(&current));
super::media::write_json_atomic(&channels_path(), &normalized)?;
let path = channels_path();
if path.exists() {
let backup = path.with_extension("json.bak");
std::fs::copy(&path, &backup).map_err(|e| format!("备份模型渠道配置失败: {e}"))?;
std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&backup)
.and_then(|file| file.sync_all())
.map_err(|e| format!("同步模型渠道备份失败: {e}"))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&backup, std::fs::Permissions::from_mode(0o600))
.map_err(|e| format!("设置模型渠道备份权限失败: {e}"))?;
}
}
super::media::write_json_atomic(&path, &normalized)?;
Ok(sanitize_doc_for_read(&normalized))
}
@@ -222,6 +314,9 @@ pub fn reveal_model_channel_key(channel_id: String) -> Result<String, String> {
.iter()
.find(|c| str_of(c, "id") == channel_id.trim())
.ok_or_else(|| format!("模型渠道不存在: {channel_id}"))?;
if normalize_api_key_ref(channel.get("apiKeyRef")).is_some() {
return Err("该渠道使用 OpenClaw SecretRef只能原样同步到 OpenClaw".into());
}
Ok(str_of(channel, "apiKey"))
}
@@ -316,4 +411,79 @@ mod tests {
let doc = normalize_channels_doc(&json!({ "channels": [entry] }), None);
assert_eq!(doc["channels"][0]["defaultModel"], "gpt-4o");
}
#[test]
fn rich_model_metadata_survives_normalization() {
let mut entry = sample_channel("ch-rich", "sk-rich");
entry["models"] = json!([{
"id": "vision",
"name": "Vision",
"input": ["text", "image"],
"reasoning": true,
"contextWindow": 200000,
"contextTokens": 160000,
"maxTokens": 8192,
"compat": { "supportsDeveloperRole": false },
"cost": { "input": 1, "output": 2 }
}]);
let doc = normalize_channels_doc(&json!({ "channels": [entry] }), None);
let model = &doc["channels"][0]["models"][0];
assert_eq!(model["input"], json!(["text", "image"]));
assert_eq!(model["reasoning"], json!(true));
assert_eq!(model["contextTokens"], json!(160000));
assert_eq!(model["maxTokens"], json!(8192));
assert_eq!(model["compat"]["supportsDeveloperRole"], json!(false));
assert_eq!(model["cost"]["output"], json!(2));
}
#[test]
fn retired_api_type_is_migrated_and_key_revision_tracks_changes() {
let mut initial = sample_channel("ch-legacy", "sk-old-same-tail");
initial["apiType"] = json!("openai-codex-responses");
let current = normalize_channels_doc(&json!({ "channels": [initial] }), None);
assert_eq!(
current["channels"][0]["apiType"],
json!("openai-chatgpt-responses")
);
assert_eq!(current["channels"][0]["credentialVersion"], json!(1));
let kept = normalize_channels_doc(
&json!({ "channels": [sample_channel("ch-legacy", "__KEEP__")] }),
Some(&current),
);
assert_eq!(kept["channels"][0]["credentialVersion"], json!(1));
let changed = normalize_channels_doc(
&json!({ "channels": [sample_channel("ch-legacy", "sk-new-same-tail")] }),
Some(&current),
);
assert_eq!(changed["channels"][0]["credentialVersion"], json!(2));
}
#[test]
fn structured_secret_ref_survives_and_plaintext_replaces_it() {
let secret_ref = json!({
"source": "env",
"provider": "default",
"id": "OPENAI_API_KEY"
});
let initial = json!({
"id": "secret-ref",
"name": "Secret Ref",
"baseUrl": "https://example.com/v1",
"apiType": "openai-responses",
"apiKey": "",
"apiKeyRef": secret_ref,
"models": [{ "id": "gpt-test" }]
});
let stored = normalize_channel(&initial, None).unwrap();
assert_eq!(stored["apiKey"], "");
assert_eq!(stored["apiKeyRef"], secret_ref);
let mut replacement = stored.clone();
replacement["apiKey"] = json!("sk-new");
let replaced = normalize_channel(&replacement, Some(&stored)).unwrap();
assert_eq!(replaced["apiKey"], "sk-new");
assert!(replaced.get("apiKeyRef").is_none());
}
}

View File

@@ -1702,6 +1702,10 @@ mod platform {
if cli.to_ascii_lowercase().ends_with(".js") {
return format!("node {} gateway", quote_batch_path(cli));
}
if cli.to_ascii_lowercase().ends_with(".cmd") || cli.to_ascii_lowercase().ends_with(".bat")
{
return format!("call {} gateway", quote_batch_path(cli));
}
format!("{} gateway", quote_batch_path(cli))
}
@@ -1709,7 +1713,7 @@ mod platform {
fs::create_dir_all(openclaw_dir).map_err(|e| format!("创建 OpenClaw 目录失败: {e}"))?;
let runner_path = openclaw_dir.join("clawpanel-gateway.cmd");
let content = format!(
"@echo off\r\ntitle {GATEWAY_WINDOW_TITLE}\r\necho Starting OpenClaw Gateway. Keep this window open after it starts.\r\necho Close this window to stop Gateway.\r\necho.\r\n{}\r\necho.\r\necho Gateway exited. You can close this window.\r\n",
"@echo off\r\ntitle {GATEWAY_WINDOW_TITLE}\r\necho Starting OpenClaw Gateway. Keep this window open after it starts.\r\necho Close this window to stop Gateway.\r\necho.\r\n{}\r\nexit /b %errorlevel%\r\n",
gateway_terminal_command(cli)
);
fs::write(&runner_path, content).map_err(|e| format!("写入 Gateway 启动脚本失败: {e}"))?;
@@ -1756,6 +1760,7 @@ mod platform {
// 外层 cmd /c 自身用 CREATE_NO_WINDOW 隐藏(短命桥接进程),
// 内部 `start` 会创建一个真正可见的新控制台窗口运行 runner.cmd。
// 使用 /C 保证 Gateway 无论以何种状态退出,包装终端都会关闭;错误由面板状态和日志展示。
let mut cmd = StdCommand::new("cmd");
cmd.args([
"/c",
@@ -1765,7 +1770,7 @@ mod platform {
openclaw_dir_str.as_str(),
"cmd",
"/D",
"/K",
"/C",
runner_path_str.as_str(),
])
.creation_flags(CREATE_NO_WINDOW)

View File

@@ -1,7 +1,7 @@
{
"$schema": "https://raw.githubusercontent.com/tauri-apps/tauri/dev/crates/tauri-config-schema/schema.json",
"productName": "ClawPanel",
"version": "0.18.5",
"version": "0.18.6",
"identifier": "ai.openclaw.clawpanel",
"build": {
"frontendDist": "../dist",