release: ship v1.1.0

This commit is contained in:
coso
2026-04-02 20:36:21 +08:00
parent 0398fc8c8e
commit 8ea5050aba
173 changed files with 9286 additions and 2275 deletions
+22
View File
@@ -235,6 +235,7 @@ pub fn run() {
.manage(automation_service_state)
.manage(commands::subagent_cmd::SubAgentSchedulerState::default())
.manage(commands::websocket_cmd::WsServiceState::default())
.manage(crate::services::companion_service::CompanionServiceState::default())
.manage(lime_gateway::telegram::TelegramGatewayState::default())
.manage(lime_gateway::discord::DiscordGatewayState::default())
.manage(lime_gateway::feishu::FeishuGatewayState::default())
@@ -444,6 +445,23 @@ pub fn run() {
});
}
{
let app_handle = app.handle().clone();
let companion_state = app
.state::<crate::services::companion_service::CompanionServiceState>()
.inner()
.clone();
tauri::async_runtime::spawn(async move {
match companion_state.start(app_handle).await {
Ok(()) => tracing::info!("[启动] Companion Pet 服务已启动"),
Err(error) => {
tracing::warn!("[启动] Companion Pet 服务启动失败: {}", error)
}
}
});
}
#[cfg(debug_assertions)]
{
let app_handle = app.handle().clone();
@@ -1129,6 +1147,10 @@ pub fn run() {
// Path utility commands
commands::config_cmd::expand_path,
commands::config_cmd::open_auth_dir,
// Companion commands
commands::companion_cmd::companion_get_pet_status,
commands::companion_cmd::companion_launch_pet,
commands::companion_cmd::companion_send_pet_command,
// OpenClaw commands
commands::openclaw_cmd::openclaw_check_installed,
commands::openclaw_cmd::openclaw_get_environment_status,
@@ -1469,6 +1469,10 @@ pub struct AgentRuntimeSpawnSubagentRequest {
#[serde(alias = "parentSessionId")]
pub parent_session_id: String,
pub message: String,
#[serde(default)]
pub name: Option<String>,
#[serde(default, alias = "teamName")]
pub team_name: Option<String>,
#[serde(default, alias = "agentType")]
pub agent_type: Option<String>,
#[serde(default)]
@@ -1499,6 +1503,8 @@ pub struct AgentRuntimeSpawnSubagentRequest {
pub system_overlay: Option<String>,
#[serde(default, alias = "outputContract")]
pub output_contract: Option<String>,
#[serde(default)]
pub cwd: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -1731,6 +1737,26 @@ mod tests {
(Utc::now() - Duration::seconds(seconds)).to_rfc3339()
}
#[test]
fn spawn_subagent_request_should_parse_current_fields() {
let request: AgentRuntimeSpawnSubagentRequest = serde_json::from_value(serde_json::json!({
"parentSessionId": "parent-1",
"message": "检查当前工具面对齐情况",
"name": "verifier",
"teamName": "delivery-team",
"agentType": "explorer",
"cwd": "/tmp/workspace"
}))
.expect("spawn subagent request should deserialize");
assert_eq!(request.parent_session_id, "parent-1");
assert_eq!(request.message, "检查当前工具面对齐情况");
assert_eq!(request.name.as_deref(), Some("verifier"));
assert_eq!(request.team_name.as_deref(), Some("delivery-team"));
assert_eq!(request.agent_type.as_deref(), Some("explorer"));
assert_eq!(request.cwd.as_deref(), Some("/tmp/workspace"));
}
#[test]
fn thread_read_should_expose_pending_request_and_waiting_incident() {
let detail = build_session_detail(
@@ -377,6 +377,7 @@ pub(crate) use session_runtime::{
resolve_session_recent_runtime_context, SessionRecentHarnessContext,
SessionRecentRuntimeContext,
};
#[allow(unused_imports)]
pub(crate) use subagent_runtime::{
agent_runtime_close_subagent_internal, agent_runtime_resume_subagent_internal,
agent_runtime_send_subagent_input_internal, agent_runtime_spawn_subagent_internal,
@@ -769,7 +769,7 @@ pub(crate) fn build_team_preference_system_prompt(
}
lines.push(
"- spawn_agent 支持这些结构化字段:blueprintRoleId、blueprintRoleLabel、teamPresetId、profileId、profileName、roleKey、skillIds、skillDirectories、theme、systemOverlay、outputContract。"
"- spawn_agent 支持这些结构化字段:name、teamName、cwd、blueprintRoleId、blueprintRoleLabel、teamPresetId、profileId、profileName、roleKey、skillIds、skillDirectories、theme、systemOverlay、outputContract。teamName 需要与 name 搭配,并依附现有 team 上下文。"
.to_string(),
);
lines.push(
@@ -9,6 +9,25 @@ const DEFAULT_WAIT_AGENT_TIMEOUT_MS: i64 = 30_000;
const MIN_WAIT_AGENT_TIMEOUT_MS: i64 = 1_000;
const MAX_WAIT_AGENT_TIMEOUT_MS: i64 = 300_000;
fn resolve_spawn_working_dir(
parent_working_dir: &std::path::Path,
requested_cwd: Option<String>,
) -> Result<std::path::PathBuf, String> {
let Some(cwd) = normalize_optional_text(requested_cwd) else {
return Ok(parent_working_dir.to_path_buf());
};
let path = std::path::PathBuf::from(&cwd);
if !path.is_absolute() {
return Err("cwd 必须是绝对路径".to_string());
}
if !path.is_dir() {
return Err(format!("cwd 不是有效目录: {cwd}"));
}
Ok(path)
}
#[derive(Debug, Clone, Serialize)]
struct SubagentStatusChangedEvent {
#[serde(rename = "type")]
@@ -153,12 +172,14 @@ fn normalize_optional_vec(values: &[String]) -> Vec<String> {
}
fn build_subagent_session_name(
explicit_name: Option<&str>,
message: &str,
agent_type: Option<&str>,
blueprint_role_label: Option<&str>,
profile_name: Option<&str>,
) -> String {
normalize_optional_text(agent_type.map(ToString::to_string))
normalize_optional_text(explicit_name.map(ToString::to_string))
.or_else(|| normalize_optional_text(agent_type.map(ToString::to_string)))
.or_else(|| normalize_optional_text(blueprint_role_label.map(ToString::to_string)))
.or_else(|| normalize_optional_text(profile_name.map(ToString::to_string)))
.or_else(|| build_subagent_task_summary(message))
@@ -169,12 +190,60 @@ fn resolve_subagent_role_hint(
request: &AgentRuntimeSpawnSubagentRequest,
customization: Option<&SubagentCustomizationState>,
) -> Option<String> {
normalize_optional_text(request.agent_type.clone())
normalize_optional_text(request.name.clone())
.or_else(|| normalize_optional_text(request.agent_type.clone()))
.or_else(|| customization.and_then(|state| state.blueprint_role_label.clone()))
.or_else(|| customization.and_then(|state| state.profile_name.clone()))
.or_else(|| customization.and_then(|state| state.role_key.clone()))
}
async fn register_spawned_teammate(
parent_session_id: &str,
child_session_id: &str,
team_name: String,
teammate_name: String,
agent_type: Option<String>,
) -> Result<(), String> {
let parent_session = read_session(parent_session_id, false, "读取父会话失败").await?;
let Some(mut team_state) = aster::session::TeamSessionState::from_session(&parent_session)
else {
return Err("当前 session 还没有 team 上下文,请先建立 team".to_string());
};
if team_state.team_name != team_name {
return Err(format!(
"team_name 不匹配:当前 team 为 {},但请求的是 {}",
team_state.team_name, team_name
));
}
if team_state.find_member_by_name(&teammate_name).is_some() {
return Err(format!("team 中已存在名为 {teammate_name} 的成员"));
}
team_state.add_or_update_member(aster::session::TeamMember::teammate(
child_session_id.to_string(),
teammate_name.clone(),
agent_type.clone(),
));
aster::session::save_team_state(parent_session_id, Some(team_state))
.await
.map_err(|error| format!("更新 team 状态失败: {error}"))?;
aster::session::save_team_membership(
child_session_id,
Some(aster::session::TeamMembershipState {
team_name,
lead_session_id: parent_session_id.to_string(),
agent_id: child_session_id.to_string(),
name: teammate_name,
agent_type,
}),
)
.await
.map_err(|error| format!("保存 team 成员信息失败: {error}"))?;
Ok(())
}
fn build_local_subagent_skill_payload(
directory: &str,
) -> Result<(SubagentSkillSummary, SubagentSkillPromptBlock), String> {
@@ -558,6 +627,11 @@ async fn create_runtime_subagent_session(
let parent_session_id =
normalize_required_text(&request.parent_session_id, "parent_session_id")?;
let message = normalize_required_text(&request.message, "message")?;
let teammate_name = normalize_optional_text(request.name.clone());
let team_name = normalize_optional_text(request.team_name.clone());
if team_name.is_some() && teammate_name.is_none() {
return Err("team_name 需要同时提供 name".to_string());
}
enforce_team_spawn_limits(&parent_session_id).await?;
let parent_session = read_session(&parent_session_id, false, "读取父会话失败").await?;
let customization = build_subagent_customization_state(request)?;
@@ -566,10 +640,13 @@ async fn create_runtime_subagent_session(
.as_ref()
.and_then(|state| state.profile_name.as_deref());
let role_hint = resolve_subagent_role_hint(request, customization.as_ref());
let working_dir =
resolve_spawn_working_dir(parent_session.working_dir.as_path(), request.cwd.clone())?;
let session = create_subagent_session(
parent_session.working_dir.clone(),
working_dir,
build_subagent_session_name(
teammate_name.as_deref(),
&message,
request.agent_type.as_deref(),
customization
@@ -614,6 +691,16 @@ async fn create_runtime_subagent_session(
"写入 subagent session metadata",
)
.await?;
if let (Some(team_name), Some(teammate_name)) = (team_name, teammate_name.clone()) {
register_spawned_teammate(
&parent_session_id,
&session.id,
team_name,
teammate_name,
normalize_optional_text(request.agent_type.clone()),
)
.await?;
}
inherit_subagent_provider(
runtime,
@@ -694,9 +781,12 @@ pub(crate) async fn agent_runtime_spawn_subagent_internal(
metadata: Some(serde_json::json!({
"subagent": {
"parent_session_id": request.parent_session_id,
"name": request.name,
"team_name": request.team_name,
"agent_type": request.agent_type,
"reasoning_effort": request.reasoning_effort,
"fork_context": request.fork_context,
"cwd": request.cwd,
"origin_tool": "spawn_agent",
"blueprint_role_id": customization.as_ref().and_then(|state| state.blueprint_role_id.clone()),
"blueprint_role_label": customization.as_ref().and_then(|state| state.blueprint_role_label.clone()),
@@ -958,3 +1048,64 @@ pub(crate) async fn agent_runtime_close_subagent_internal(
changed_session_ids: changed_ids,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_build_subagent_session_name_prefers_explicit_name() {
let name = build_subagent_session_name(
Some("verifier"),
"检查当前 team runtime 差异",
Some("explorer"),
Some("分析"),
Some("代码分析员"),
);
assert_eq!(name, "verifier");
}
#[test]
fn test_resolve_subagent_role_hint_prefers_explicit_name() {
let request = AgentRuntimeSpawnSubagentRequest {
parent_session_id: "parent-1".to_string(),
message: "定位当前 team runtime 差异".to_string(),
name: Some("verifier".to_string()),
team_name: Some("delivery-team".to_string()),
agent_type: Some("explorer".to_string()),
model: None,
reasoning_effort: None,
fork_context: false,
blueprint_role_id: None,
blueprint_role_label: Some("分析".to_string()),
profile_id: None,
profile_name: Some("代码分析员".to_string()),
role_key: Some("explorer".to_string()),
skill_ids: Vec::new(),
skill_directories: Vec::new(),
team_preset_id: None,
theme: None,
system_overlay: None,
output_contract: None,
cwd: None,
};
assert_eq!(
resolve_subagent_role_hint(&request, None).as_deref(),
Some("verifier")
);
}
#[test]
fn test_resolve_spawn_working_dir_uses_requested_absolute_directory() {
let parent = tempfile::tempdir().expect("parent tempdir");
let child = tempfile::tempdir().expect("child tempdir");
let resolved =
resolve_spawn_working_dir(parent.path(), Some(child.path().display().to_string()))
.expect("cwd override should resolve");
assert_eq!(resolved, child.path());
}
}
@@ -2604,6 +2604,8 @@ mod tests {
let customization = build_subagent_customization_state(&AgentRuntimeSpawnSubagentRequest {
parent_session_id: "parent-1".to_string(),
message: "定位当前 team runtime 差异".to_string(),
name: None,
team_name: None,
agent_type: Some("Image #1".to_string()),
model: None,
reasoning_effort: None,
@@ -2619,6 +2621,7 @@ mod tests {
theme: None,
system_overlay: None,
output_contract: None,
cwd: None,
})
.expect("build customization state")
.expect("customization should exist");
@@ -408,6 +408,8 @@ impl Tool for SubAgentTaskTool {
AgentRuntimeSpawnSubagentRequest {
parent_session_id,
message: build_subagent_task_runtime_message(&input, &task, role),
name: None,
team_name: None,
agent_type: Some(role.to_string()),
model: input.model.clone(),
reasoning_effort: None,
@@ -423,6 +425,7 @@ impl Tool for SubAgentTaskTool {
theme: None,
system_overlay: None,
output_contract: None,
cwd: None,
},
)
.await
@@ -517,10 +520,7 @@ fn build_agent_control_tool_config(
runtime: SubagentControlRuntime,
) -> aster::tools::AgentControlToolConfig {
let spawn_runtime = runtime.clone();
let send_runtime = runtime.clone();
let wait_runtime = runtime.clone();
let resume_runtime = runtime.clone();
let close_runtime = runtime;
let send_runtime = runtime;
aster::tools::AgentControlToolConfig::new()
.with_spawn_agent_callback(Arc::new(move |request| {
@@ -531,6 +531,8 @@ fn build_agent_control_tool_config(
AgentRuntimeSpawnSubagentRequest {
parent_session_id: request.parent_session_id,
message: request.message,
name: request.name,
team_name: request.team_name,
agent_type: request.agent_type,
model: request.model,
reasoning_effort: request.reasoning_effort,
@@ -546,6 +548,7 @@ fn build_agent_control_tool_config(
theme: request.theme,
system_overlay: request.system_overlay,
output_contract: request.output_contract,
cwd: request.cwd,
},
)
.await?;
@@ -576,80 +579,6 @@ fn build_agent_control_tool_config(
})
})
}))
.with_wait_agent_callback(Arc::new(move |request| {
let runtime = wait_runtime.clone();
Box::pin(async move {
let response = agent_runtime_wait_subagents_internal(
&runtime,
AgentRuntimeWaitSubagentsRequest {
ids: request.ids,
timeout_ms: request.timeout_ms,
},
)
.await?;
let status = response
.status
.into_iter()
.map(|(id, status)| {
serde_json::to_value(status)
.map(|value| (id, value))
.map_err(|error| format!("wait_agent 状态序列化失败: {error}"))
})
.collect::<Result<std::collections::BTreeMap<_, _>, _>>()?;
Ok(aster::tools::WaitAgentResponse {
status,
timed_out: response.timed_out,
extra: std::collections::BTreeMap::new(),
})
})
}))
.with_resume_agent_callback(Arc::new(move |request| {
let runtime = resume_runtime.clone();
Box::pin(async move {
let response = agent_runtime_resume_subagent_internal(
&runtime,
AgentRuntimeResumeSubagentRequest { id: request.id },
)
.await?;
let mut extra = std::collections::BTreeMap::new();
extra.insert(
"cascade_session_ids".to_string(),
serde_json::to_value(response.cascade_session_ids)
.map_err(|error| format!("resume_agent 级联会话序列化失败: {error}"))?,
);
Ok(aster::tools::ResumeAgentResponse {
status: serde_json::to_value(response.status)
.map_err(|error| format!("resume_agent 状态序列化失败: {error}"))?,
changed_session_ids: response.changed_session_ids,
extra,
})
})
}))
.with_close_agent_callback(Arc::new(move |request| {
let runtime = close_runtime.clone();
Box::pin(async move {
let response = agent_runtime_close_subagent_internal(
&runtime,
AgentRuntimeCloseSubagentRequest { id: request.id },
)
.await?;
let mut extra = std::collections::BTreeMap::new();
extra.insert(
"cascade_session_ids".to_string(),
serde_json::to_value(response.cascade_session_ids)
.map_err(|error| format!("close_agent 级联会话序列化失败: {error}"))?,
);
Ok(aster::tools::CloseAgentResponse {
previous_status: serde_json::to_value(response.previous_status)
.map_err(|error| format!("close_agent 状态序列化失败: {error}"))?,
changed_session_ids: response.changed_session_ids,
extra,
})
})
}))
}
pub(super) fn register_subagent_runtime_tools(
+28
View File
@@ -0,0 +1,28 @@
use crate::services::companion_service::{
self, CompanionLaunchPetRequest, CompanionLaunchPetResult, CompanionPetCommandRequest,
CompanionPetSendResult, CompanionPetStatus, CompanionServiceState,
};
use tauri::State;
#[tauri::command]
pub async fn companion_get_pet_status(
companion_state: State<'_, CompanionServiceState>,
) -> Result<CompanionPetStatus, String> {
companion_service::get_pet_status_global(companion_state.inner()).await
}
#[tauri::command]
pub async fn companion_launch_pet(
companion_state: State<'_, CompanionServiceState>,
request: Option<CompanionLaunchPetRequest>,
) -> Result<CompanionLaunchPetResult, String> {
companion_service::launch_pet_global(companion_state.inner(), request.unwrap_or_default()).await
}
#[tauri::command]
pub async fn companion_send_pet_command(
companion_state: State<'_, CompanionServiceState>,
request: CompanionPetCommandRequest,
) -> Result<CompanionPetSendResult, String> {
companion_service::send_pet_command_global(companion_state.inner(), request).await
}
+1
View File
@@ -11,6 +11,7 @@ pub mod browser_profile_cmd;
pub mod browser_runtime_cmd;
pub mod channels_cmd;
pub mod claw_solution_cmd;
pub mod companion_cmd;
pub mod config_cmd;
pub mod connect_cmd;
pub mod connection_cmd;
+25 -19
View File
@@ -5,9 +5,9 @@
use crate::models::model_registry::{
EnhancedModelMetadata, ModelSyncState, ModelTier, ProviderAliasConfig, UserModelPreference,
};
use lime_server_utils::load_model_registry_provider_ids_from_resources;
use lime_services::model_registry_service::{FetchModelsResult, ModelRegistryService};
use serde::Serialize;
use std::collections::BTreeSet;
use std::sync::Arc;
use tauri::State;
use tokio::sync::RwLock;
@@ -36,20 +36,12 @@ pub async fn get_model_registry(
/// 获取模型注册表中所有 provider_id(去重且有序)
///
/// provider_id 来源于 `src-tauri/resources/models/providers/*.json` 加载结果。
/// provider_id 以 `src-tauri/resources/models/index.json` 的 providers 列表为唯一真相源。
#[tauri::command]
pub async fn get_model_registry_provider_ids(
state: State<'_, ModelRegistryState>,
_state: State<'_, ModelRegistryState>,
) -> Result<Vec<String>, String> {
let guard = state.read().await;
let service = guard
.as_ref()
.ok_or_else(|| "模型注册服务未初始化".to_string())?;
let models = service.get_all_models().await;
let provider_ids: BTreeSet<String> = models.into_iter().map(|m| m.provider_id).collect();
Ok(provider_ids.into_iter().collect())
load_model_registry_provider_ids_from_resources()
}
/// 搜索模型
@@ -272,18 +264,32 @@ pub async fn fetch_provider_models_auto(
.get_provider(&db, &provider_id)?
.ok_or_else(|| format!("Provider 不存在: {provider_id}"))?;
// 获取 API Key
let api_key = api_key_service
.0
.get_next_api_key(&db, &provider_id)?
.ok_or_else(|| format!("Provider {provider_id} 没有可用的 API Key"))?;
// 获取 API Host
let api_host = provider.provider.api_host.clone();
if api_host.is_empty() {
return Err("Provider 没有配置 API Host".to_string());
}
let provider_type = provider.provider.provider_type;
let requires_api_key = ModelRegistryService::requires_api_key_for_model_fetch(
&provider_id,
&api_host,
provider_type,
);
// 获取 API Key(支持本地免 Key 渠道)
let api_key = if requires_api_key {
api_key_service
.0
.get_next_api_key(&db, &provider_id)?
.ok_or_else(|| format!("Provider {provider_id} 没有可用的 API Key"))?
} else {
api_key_service
.0
.get_next_api_key(&db, &provider_id)?
.unwrap_or_default()
};
// 调用模型注册服务
let guard = state.read().await;
let service = guard
@@ -295,7 +301,7 @@ pub async fn fetch_provider_models_auto(
&provider_id,
&api_host,
&api_key,
Some(provider.provider.provider_type),
Some(provider_type),
&provider.provider.custom_models,
)
.await
+17 -2
View File
@@ -3684,6 +3684,10 @@ mod tests {
use crate::services::browser_profile_service::sanitize_browser_profile_key;
use rusqlite::Connection;
use std::sync::{Arc, Mutex};
use tokio::sync::Mutex as AsyncMutex;
static BROWSER_RUNTIME_AUDIT_TEST_LOCK: Lazy<AsyncMutex<()>> =
Lazy::new(|| AsyncMutex::new(()));
fn setup_db() -> DbConnection {
let conn = Connection::open_in_memory().unwrap();
@@ -3995,6 +3999,7 @@ mod tests {
#[tokio::test]
async fn browser_runtime_audit_should_store_launch_metadata() {
let _guard = BROWSER_RUNTIME_AUDIT_TEST_LOCK.lock().await;
BROWSER_RUNTIME_AUDIT_LOGS.lock().await.clear();
append_browser_runtime_launch_audit(BrowserRuntimeLaunchAuditInput {
@@ -4018,7 +4023,10 @@ mod tests {
let logs = get_browser_action_audit_logs(Some(5))
.await
.expect("audit logs should be readable");
let record = logs.first().expect("launch audit must exist");
let record = logs
.iter()
.find(|record| matches!(record.kind, BrowserRuntimeAuditKind::Launch))
.expect("launch audit must exist");
assert!(matches!(record.kind, BrowserRuntimeAuditKind::Launch));
assert_eq!(
record.profile_key.as_deref(),
@@ -4041,6 +4049,7 @@ mod tests {
#[tokio::test]
async fn browser_runtime_audit_should_store_action_session_keys() {
let _guard = BROWSER_RUNTIME_AUDIT_TEST_LOCK.lock().await;
BROWSER_RUNTIME_AUDIT_LOGS.lock().await.clear();
append_browser_runtime_audit(BrowserRuntimeAuditRecord::action(
@@ -4064,7 +4073,13 @@ mod tests {
let logs = get_browser_action_audit_logs(Some(5))
.await
.expect("audit logs should be readable");
let record = logs.first().expect("action audit must exist");
let record = logs
.iter()
.find(|record| {
matches!(record.kind, BrowserRuntimeAuditKind::Action)
&& record.id == "browser-action-1"
})
.expect("action audit must exist");
assert!(matches!(record.kind, BrowserRuntimeAuditKind::Action));
assert_eq!(record.id, "browser-action-1");
assert_eq!(record.session_id.as_deref(), Some("session-42"));
+10
View File
@@ -4,7 +4,9 @@
mod agent_sessions;
mod app_runtime;
mod automation;
mod browser;
mod companion;
mod content;
mod logs;
mod memory;
@@ -90,6 +92,14 @@ pub async fn handle_command(
return Ok(result);
}
if let Some(result) = companion::try_handle(state, cmd, args.as_ref()).await? {
return Ok(result);
}
if let Some(result) = automation::try_handle(state, cmd, args.as_ref()).await? {
return Ok(result);
}
if let Some(result) = logs::try_handle(state, cmd, args.as_ref()).await? {
return Ok(result);
}
@@ -0,0 +1,158 @@
use super::{
args_or_default, get_string_arg, parse_nested_arg, parse_optional_nested_arg,
require_app_handle,
};
use crate::dev_bridge::DevBridgeState;
use serde_json::Value as JsonValue;
use tauri::Manager;
type DynError = Box<dyn std::error::Error>;
pub(super) async fn try_handle(
state: &DevBridgeState,
cmd: &str,
args: Option<&JsonValue>,
) -> Result<Option<JsonValue>, DynError> {
if !matches!(
cmd,
"get_automation_scheduler_config"
| "update_automation_scheduler_config"
| "get_automation_status"
| "get_automation_jobs"
| "get_automation_job"
| "create_automation_job"
| "update_automation_job"
| "delete_automation_job"
| "run_automation_job_now"
| "get_automation_health"
| "get_automation_run_history"
| "preview_automation_schedule"
| "validate_automation_schedule"
) {
return Ok(None);
}
let app_handle = require_app_handle(state)?;
let app_state = app_handle.state::<crate::app::AppState>();
let automation_state =
app_handle.state::<crate::services::automation_service::AutomationServiceState>();
let result = match cmd {
"get_automation_scheduler_config" => serde_json::to_value(
crate::commands::automation_cmd::get_automation_scheduler_config(app_state).await?,
)?,
"update_automation_scheduler_config" => {
let config = parse_nested_arg::<
crate::commands::automation_cmd::AutomationSchedulerConfigResponse,
>(&args_or_default(args), "config")?;
crate::commands::automation_cmd::update_automation_scheduler_config(
app_state,
automation_state,
config,
app_handle.clone(),
)
.await?;
JsonValue::Null
}
"get_automation_status" => serde_json::to_value(
crate::commands::automation_cmd::get_automation_status(automation_state).await?,
)?,
"get_automation_jobs" => serde_json::to_value(
crate::commands::automation_cmd::get_automation_jobs(automation_state).await?,
)?,
"get_automation_job" => {
let args = args_or_default(args);
let id = get_string_arg(&args, "id", "id")?;
serde_json::to_value(
crate::commands::automation_cmd::get_automation_job(automation_state, id).await?,
)?
}
"create_automation_job" => {
let request = parse_nested_arg::<crate::commands::automation_cmd::AutomationJobRequest>(
&args_or_default(args),
"request",
)?;
serde_json::to_value(
crate::commands::automation_cmd::create_automation_job(automation_state, request)
.await?,
)?
}
"update_automation_job" => {
let args = args_or_default(args);
let id = get_string_arg(&args, "id", "id")?;
let request = parse_nested_arg::<
crate::commands::automation_cmd::UpdateAutomationJobRequest,
>(&args, "request")?;
serde_json::to_value(
crate::commands::automation_cmd::update_automation_job(
automation_state,
id,
request,
)
.await?,
)?
}
"delete_automation_job" => {
let args = args_or_default(args);
let id = get_string_arg(&args, "id", "id")?;
serde_json::to_value(
crate::commands::automation_cmd::delete_automation_job(automation_state, id)
.await?,
)?
}
"run_automation_job_now" => {
let args = args_or_default(args);
let id = get_string_arg(&args, "id", "id")?;
serde_json::to_value(
crate::commands::automation_cmd::run_automation_job_now(automation_state, id)
.await?,
)?
}
"get_automation_health" => {
let query = parse_optional_nested_arg::<
crate::services::automation_service::health::AutomationHealthQuery,
>(&args_or_default(args), "query")?;
serde_json::to_value(
crate::commands::automation_cmd::get_automation_health(automation_state, query)
.await?,
)?
}
"get_automation_run_history" => {
let args = args_or_default(args);
let id = get_string_arg(&args, "id", "id")?;
let limit = args
.get("limit")
.and_then(|value| value.as_u64())
.map(|value| value as usize);
serde_json::to_value(
crate::commands::automation_cmd::get_automation_run_history(
automation_state,
id,
limit,
)
.await?,
)?
}
"preview_automation_schedule" => {
let schedule = parse_nested_arg::<lime_core::config::TaskSchedule>(
&args_or_default(args),
"schedule",
)?;
serde_json::to_value(
crate::commands::automation_cmd::preview_automation_schedule(schedule).await?,
)?
}
"validate_automation_schedule" => {
let schedule = parse_nested_arg::<lime_core::config::TaskSchedule>(
&args_or_default(args),
"schedule",
)?;
serde_json::to_value(
crate::commands::automation_cmd::validate_automation_schedule(schedule).await?,
)?
}
_ => unreachable!("已通过前置 matches! 过滤 automation 命令"),
};
Ok(Some(result))
}
@@ -0,0 +1,51 @@
use super::{args_or_default, parse_nested_arg, parse_optional_nested_arg, require_app_handle};
use crate::dev_bridge::DevBridgeState;
use serde_json::Value as JsonValue;
use tauri::Manager;
type DynError = Box<dyn std::error::Error>;
pub(super) async fn try_handle(
state: &DevBridgeState,
cmd: &str,
args: Option<&JsonValue>,
) -> Result<Option<JsonValue>, DynError> {
let app_handle = match cmd {
"companion_get_pet_status" | "companion_launch_pet" | "companion_send_pet_command" => {
require_app_handle(state)?
}
_ => return Ok(None),
};
let companion_state =
app_handle.state::<crate::services::companion_service::CompanionServiceState>();
let result = match cmd {
"companion_get_pet_status" => serde_json::to_value(
crate::services::companion_service::get_pet_status_global(&companion_state).await?,
)?,
"companion_launch_pet" => {
let request = parse_optional_nested_arg::<
crate::services::companion_service::CompanionLaunchPetRequest,
>(&args_or_default(args), "request")?
.unwrap_or_default();
serde_json::to_value(
crate::services::companion_service::launch_pet_global(&companion_state, request)
.await?,
)?
}
"companion_send_pet_command" => {
let request: crate::services::companion_service::CompanionPetCommandRequest =
parse_nested_arg(&args_or_default(args), "request")?;
serde_json::to_value(
crate::services::companion_service::send_pet_command_global(
&companion_state,
request,
)
.await?,
)?
}
_ => unreachable!("已通过前置判断过滤 companion 命令"),
};
Ok(Some(result))
}
+2 -35
View File
@@ -4,32 +4,11 @@ use serde_json::Value as JsonValue;
type DynError = Box<dyn std::error::Error>;
fn load_model_registry_provider_ids_from_db(
state: &DevBridgeState,
) -> Result<Vec<String>, DynError> {
let Some(db) = &state.db else {
return Ok(vec![]);
};
let conn = db.lock().map_err(|e| format!("数据库锁定失败: {e}"))?;
let mut stmt = conn.prepare(
"SELECT DISTINCT provider_id FROM model_registry WHERE provider_id IS NOT NULL ORDER BY provider_id",
)?;
let rows = stmt.query_map([], |row| row.get::<_, String>(0))?;
let mut provider_ids = Vec::new();
for row in rows {
provider_ids.push(row?);
}
Ok(provider_ids)
}
pub(super) async fn try_handle(
state: &DevBridgeState,
cmd: &str,
) -> Result<Option<JsonValue>, DynError> {
let result = match cmd {
let result: JsonValue = match cmd {
"get_models" => serde_json::json!({
"data": [
{"id": "claude-sonnet-4-20250514", "object": "model", "owned_by": "anthropic"},
@@ -75,19 +54,7 @@ pub(super) async fn try_handle(
serde_json::json!(service.force_reload().await?)
}
"get_model_registry_provider_ids" => {
match load_model_registry_provider_ids_from_resources() {
Ok(provider_ids) => serde_json::to_value(provider_ids)?,
Err(resource_error) => {
let fallback = load_model_registry_provider_ids_from_db(state)?;
if fallback.is_empty() {
return Err(format!(
"获取模型 Provider ID 失败(resources 与数据库均不可用): {resource_error}"
)
.into());
}
serde_json::to_value(fallback)?
}
}
serde_json::to_value(load_model_registry_provider_ids_from_resources()?)?
}
_ => return Ok(None),
};
+782
View File
@@ -0,0 +1,782 @@
use axum::{
extract::{
ws::{Message, WebSocket, WebSocketUpgrade},
State,
},
response::IntoResponse,
routing::get,
Router,
};
use futures::{SinkExt, StreamExt};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use std::{path::PathBuf, process::Command, sync::Arc};
use tauri::{AppHandle, Emitter, Manager};
use tokio::sync::{mpsc, Mutex, RwLock};
use uuid::Uuid;
const DEFAULT_COMPANION_HOST: &str = "127.0.0.1";
const DEFAULT_COMPANION_PORT: u16 = 45554;
const DEFAULT_COMPANION_PATH: &str = "/companion/pet";
const DEFAULT_CLIENT_ID: &str = "lime";
const DEFAULT_PROTOCOL_VERSION: u32 = 1;
const COMPANION_ENV_APP_PATH: &str = "LIME_PET_APP_PATH";
#[cfg(target_os = "macos")]
const MACOS_PET_APP_NAME: &str = "Lime Pet";
#[cfg(target_os = "windows")]
const WINDOWS_PET_EXE_NAME: &str = "Lime Pet.exe";
pub const COMPANION_PET_STATUS_EVENT: &str = "companion-pet-status";
pub const COMPANION_OPEN_PROVIDER_SETTINGS_EVENT: &str = "companion-open-provider-settings";
fn default_companion_endpoint() -> String {
format!("ws://{DEFAULT_COMPANION_HOST}:{DEFAULT_COMPANION_PORT}{DEFAULT_COMPANION_PATH}")
}
fn empty_payload() -> Value {
Value::Object(Map::new())
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum CompanionPetVisualState {
Hidden,
Idle,
Walking,
Thinking,
Done,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompanionPetStatus {
pub endpoint: String,
pub server_listening: bool,
pub connected: bool,
pub client_id: Option<String>,
pub platform: Option<String>,
pub capabilities: Vec<String>,
pub last_event: Option<String>,
pub last_error: Option<String>,
pub last_state: Option<CompanionPetVisualState>,
}
impl Default for CompanionPetStatus {
fn default() -> Self {
Self {
endpoint: default_companion_endpoint(),
server_listening: false,
connected: false,
client_id: None,
platform: None,
capabilities: Vec::new(),
last_event: None,
last_error: None,
last_state: None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct CompanionLaunchPetRequest {
pub app_path: Option<String>,
pub endpoint: Option<String>,
pub client_id: Option<String>,
pub protocol_version: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompanionLaunchPetResult {
pub launched: bool,
pub resolved_path: Option<String>,
pub endpoint: String,
pub message: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompanionPetCommandRequest {
pub event: String,
#[serde(default = "empty_payload")]
pub payload: Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompanionPetSendResult {
pub delivered: bool,
pub connected: bool,
}
#[derive(Debug, Clone, Serialize)]
struct CompanionEnvelope {
protocol_version: u32,
event: String,
payload: Value,
}
#[derive(Debug, Clone, Deserialize)]
struct CompanionIncomingEnvelope {
protocol_version: u32,
event: String,
#[serde(default = "empty_payload")]
payload: Value,
}
#[derive(Debug, Clone, Deserialize, Default)]
struct CompanionReadyPayload {
client_id: Option<String>,
platform: Option<String>,
#[serde(default)]
capabilities: Vec<String>,
}
#[derive(Debug, Clone)]
struct ActivePetSender {
connection_id: String,
tx: mpsc::UnboundedSender<String>,
}
#[derive(Debug, Clone, Default)]
struct CompanionRuntime {
status: CompanionPetStatus,
active_connection_id: Option<String>,
}
#[derive(Clone)]
struct CompanionRouterState {
app_handle: AppHandle,
service: CompanionServiceState,
}
#[derive(Clone)]
pub struct CompanionServiceState {
app_handle: Arc<RwLock<Option<AppHandle>>>,
runtime: Arc<RwLock<CompanionRuntime>>,
sender: Arc<Mutex<Option<ActivePetSender>>>,
start_lock: Arc<Mutex<()>>,
}
impl Default for CompanionServiceState {
fn default() -> Self {
Self {
app_handle: Arc::new(RwLock::new(None)),
runtime: Arc::new(RwLock::new(CompanionRuntime {
status: CompanionPetStatus::default(),
active_connection_id: None,
})),
sender: Arc::new(Mutex::new(None)),
start_lock: Arc::new(Mutex::new(())),
}
}
}
impl CompanionServiceState {
pub async fn start(&self, app_handle: AppHandle) -> Result<(), String> {
self.set_app_handle(app_handle.clone()).await;
let _guard = self.start_lock.lock().await;
if self.snapshot().await.server_listening {
return Ok(());
}
let bind_addr = format!("{DEFAULT_COMPANION_HOST}:{DEFAULT_COMPANION_PORT}");
let listener = match tokio::net::TcpListener::bind(&bind_addr).await {
Ok(listener) => listener,
Err(error) => {
self.update_runtime(|runtime| {
runtime.status.server_listening = false;
runtime.status.last_error = Some(format!("Companion 服务监听失败: {error}"));
})
.await;
return Err(format!("Companion 服务监听失败: {error}"));
}
};
self.update_runtime(|runtime| {
runtime.status.server_listening = true;
runtime.status.last_error = None;
runtime.status.endpoint = default_companion_endpoint();
})
.await;
let service = self.clone();
let router = Router::new()
.route(DEFAULT_COMPANION_PATH, get(companion_pet_ws))
.with_state(CompanionRouterState {
app_handle,
service: self.clone(),
});
tokio::spawn(async move {
if let Err(error) = axum::serve(listener, router).await {
tracing::error!("[Companion] 服务运行失败: {}", error);
service
.mark_server_stopped(Some(format!("Companion 服务运行失败: {error}")))
.await;
}
});
tracing::info!(
"[Companion] 已监听桌宠入口: {}",
default_companion_endpoint()
);
Ok(())
}
pub async fn snapshot(&self) -> CompanionPetStatus {
self.runtime.read().await.status.clone()
}
pub async fn launch_pet(
&self,
request: CompanionLaunchPetRequest,
) -> Result<CompanionLaunchPetResult, String> {
let endpoint = request
.endpoint
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(default_companion_endpoint);
let client_id = request
.client_id
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| DEFAULT_CLIENT_ID.to_string());
let protocol_version = request.protocol_version.unwrap_or(DEFAULT_PROTOCOL_VERSION);
let Some(target) = resolve_launch_target(request.app_path.as_deref()) else {
return Ok(CompanionLaunchPetResult {
launched: false,
resolved_path: None,
endpoint,
message: Some(
"未找到 Lime Pet 可执行产物,请先安装桌宠应用或通过 app_path 显式指定。"
.to_string(),
),
});
};
let mut command = Command::new(&target.exec_path);
command
.arg("--connect")
.arg(&endpoint)
.arg("--client-id")
.arg(&client_id)
.arg("--protocol-version")
.arg(protocol_version.to_string());
match command.spawn() {
Ok(_) => Ok(CompanionLaunchPetResult {
launched: true,
resolved_path: Some(target.exec_path.to_string_lossy().to_string()),
endpoint,
message: None,
}),
Err(error) => Ok(CompanionLaunchPetResult {
launched: false,
resolved_path: Some(target.exec_path.to_string_lossy().to_string()),
endpoint,
message: Some(format!("启动 Lime Pet 失败: {error}")),
}),
}
}
pub async fn send_pet_command(
&self,
request: CompanionPetCommandRequest,
) -> Result<CompanionPetSendResult, String> {
let payload = CompanionEnvelope {
protocol_version: DEFAULT_PROTOCOL_VERSION,
event: request.event.clone(),
payload: request.payload.clone(),
};
let serialized = serde_json::to_string(&payload).map_err(|e| e.to_string())?;
let active_sender = self.sender.lock().await.clone();
let Some(active_sender) = active_sender else {
return Ok(CompanionPetSendResult {
delivered: false,
connected: false,
});
};
if active_sender.tx.send(serialized).is_err() {
self.mark_connection_closed(
&active_sender.connection_id,
"桌宠连接不可用,命令未送达".to_string(),
)
.await;
return Ok(CompanionPetSendResult {
delivered: false,
connected: false,
});
}
self.record_outbound_command(&request.event, &request.payload)
.await;
Ok(CompanionPetSendResult {
delivered: true,
connected: true,
})
}
async fn set_app_handle(&self, app_handle: AppHandle) {
let mut guard = self.app_handle.write().await;
*guard = Some(app_handle);
}
async fn emit_status(&self, status: CompanionPetStatus) {
let app_handle = self.app_handle.read().await.clone();
if let Some(app_handle) = app_handle {
if let Err(error) = app_handle.emit(COMPANION_PET_STATUS_EVENT, &status) {
tracing::warn!("[Companion] 发送状态事件失败: {}", error);
}
}
}
async fn update_runtime<F>(&self, mutate: F)
where
F: FnOnce(&mut CompanionRuntime),
{
let snapshot = {
let mut runtime = self.runtime.write().await;
mutate(&mut runtime);
runtime.status.clone()
};
self.emit_status(snapshot).await;
}
async fn mark_server_stopped(&self, reason: Option<String>) {
let mut sender = self.sender.lock().await;
sender.take();
drop(sender);
self.update_runtime(|runtime| {
runtime.status.server_listening = false;
runtime.status.connected = false;
runtime.status.client_id = None;
runtime.status.platform = None;
runtime.status.capabilities.clear();
runtime.status.last_error = reason;
runtime.active_connection_id = None;
})
.await;
}
async fn attach_sender(&self, connection_id: String, tx: mpsc::UnboundedSender<String>) {
let mut guard = self.sender.lock().await;
*guard = Some(ActivePetSender { connection_id, tx });
}
async fn mark_connection_open(&self, connection_id: String) {
self.update_runtime(|runtime| {
runtime.active_connection_id = Some(connection_id);
runtime.status.connected = true;
runtime.status.client_id = None;
runtime.status.platform = None;
runtime.status.capabilities.clear();
runtime.status.last_event = Some("pet.connected".to_string());
runtime.status.last_error = None;
})
.await;
}
async fn mark_connection_ready(&self, connection_id: &str, payload: CompanionReadyPayload) {
self.update_runtime(|runtime| {
if runtime.active_connection_id.as_deref() != Some(connection_id) {
return;
}
runtime.status.connected = true;
runtime.status.client_id = payload.client_id;
runtime.status.platform = payload.platform;
runtime.status.capabilities = payload.capabilities;
runtime.status.last_event = Some("pet.ready".to_string());
runtime.status.last_error = None;
})
.await;
}
async fn mark_connection_closed(&self, connection_id: &str, reason: String) {
{
let mut sender = self.sender.lock().await;
if sender.as_ref().map(|active| active.connection_id.as_str()) == Some(connection_id) {
sender.take();
}
}
self.update_runtime(|runtime| {
if runtime.active_connection_id.as_deref() != Some(connection_id) {
return;
}
runtime.active_connection_id = None;
runtime.status.connected = false;
runtime.status.client_id = None;
runtime.status.platform = None;
runtime.status.capabilities.clear();
runtime.status.last_event = Some("pet.disconnected".to_string());
runtime.status.last_error = Some(reason);
})
.await;
}
async fn handle_incoming_message(
&self,
app_handle: &AppHandle,
connection_id: &str,
text: &str,
) {
let envelope = match serde_json::from_str::<CompanionIncomingEnvelope>(text) {
Ok(envelope) => envelope,
Err(error) => {
tracing::warn!("[Companion] 忽略无法解析的桌宠消息: {}", error);
return;
}
};
if envelope.protocol_version != DEFAULT_PROTOCOL_VERSION {
self.update_runtime(|runtime| {
if runtime.active_connection_id.as_deref() != Some(connection_id) {
return;
}
runtime.status.last_error =
Some(format!("桌宠协议版本不兼容: {}", envelope.protocol_version));
})
.await;
return;
}
let should_focus_main_window = matches!(
envelope.event.as_str(),
"pet.clicked" | "pet.open_chat" | "pet.open_provider_settings"
);
self.update_runtime(|runtime| {
if runtime.active_connection_id.as_deref() != Some(connection_id) {
return;
}
runtime.status.last_event = Some(envelope.event.clone());
runtime.status.last_error = None;
})
.await;
if envelope.event == "pet.ready" {
match serde_json::from_value::<CompanionReadyPayload>(envelope.payload.clone()) {
Ok(payload) => {
self.mark_connection_ready(connection_id, payload).await;
}
Err(error) => {
self.update_runtime(|runtime| {
if runtime.active_connection_id.as_deref() != Some(connection_id) {
return;
}
runtime.status.last_error =
Some(format!("桌宠 ready 负载解析失败: {error}"));
})
.await;
}
}
}
if should_focus_main_window {
reveal_main_window(app_handle);
}
if envelope.event == "pet.open_provider_settings" {
if let Err(error) = app_handle.emit(COMPANION_OPEN_PROVIDER_SETTINGS_EVENT, ()) {
tracing::warn!("[Companion] 发送服务商设置跳转事件失败: {}", error);
}
}
}
async fn record_outbound_command(&self, event: &str, payload: &Value) {
let event_name = event.to_string();
let next_state = match event {
"pet.hide" => Some(CompanionPetVisualState::Hidden),
"pet.show" => Some(CompanionPetVisualState::Walking),
"pet.state_changed" => payload
.get("state")
.and_then(Value::as_str)
.and_then(parse_visual_state),
_ => None,
};
self.update_runtime(|runtime| {
runtime.status.last_event = Some(event_name);
runtime.status.last_error = None;
if let Some(next_state) = next_state {
runtime.status.last_state = Some(next_state);
}
})
.await;
}
}
fn parse_visual_state(value: &str) -> Option<CompanionPetVisualState> {
match value.trim() {
"hidden" => Some(CompanionPetVisualState::Hidden),
"idle" => Some(CompanionPetVisualState::Idle),
"walking" => Some(CompanionPetVisualState::Walking),
"thinking" => Some(CompanionPetVisualState::Thinking),
"done" => Some(CompanionPetVisualState::Done),
_ => None,
}
}
fn reveal_main_window(app_handle: &AppHandle) {
let Some(window) = app_handle.get_webview_window("main") else {
tracing::warn!("[Companion] 未找到主窗口,无法响应桌宠点击");
return;
};
if let Err(error) = window.unminimize() {
tracing::warn!("[Companion] 主窗口取消最小化失败: {}", error);
}
if let Err(error) = window.show() {
tracing::warn!("[Companion] 主窗口显示失败: {}", error);
}
if let Err(error) = window.set_focus() {
tracing::warn!("[Companion] 主窗口聚焦失败: {}", error);
}
}
#[derive(Debug)]
struct LaunchTarget {
exec_path: PathBuf,
}
fn resolve_launch_target(explicit_path: Option<&str>) -> Option<LaunchTarget> {
let mut candidates = Vec::new();
if let Some(explicit_path) = explicit_path
.map(str::trim)
.filter(|value| !value.is_empty())
{
candidates.push(expand_user_path(explicit_path));
}
if let Ok(env_path) = std::env::var(COMPANION_ENV_APP_PATH) {
let env_path = env_path.trim();
if !env_path.is_empty() {
candidates.push(expand_user_path(env_path));
}
}
candidates.extend(default_launch_candidates());
candidates.into_iter().find_map(normalize_launch_target)
}
fn expand_user_path(value: &str) -> PathBuf {
if value == "~" {
return dirs::home_dir().unwrap_or_else(|| PathBuf::from(value));
}
if let Some(suffix) = value.strip_prefix("~/") {
if let Some(home_dir) = dirs::home_dir() {
return home_dir.join(suffix);
}
}
PathBuf::from(value)
}
#[cfg(target_os = "macos")]
fn default_launch_candidates() -> Vec<PathBuf> {
let mut candidates = Vec::new();
if let Some(home_dir) = dirs::home_dir() {
candidates.push(home_dir.join("Applications/Lime Pet.app"));
}
candidates.push(PathBuf::from("/Applications/Lime Pet.app"));
candidates
}
#[cfg(target_os = "windows")]
fn default_launch_candidates() -> Vec<PathBuf> {
let mut candidates = Vec::new();
if let Some(local_dir) = dirs::data_local_dir() {
candidates.push(
local_dir
.join("Programs")
.join("Lime Pet")
.join(WINDOWS_PET_EXE_NAME),
);
}
candidates
}
#[cfg(not(any(target_os = "macos", target_os = "windows")))]
fn default_launch_candidates() -> Vec<PathBuf> {
Vec::new()
}
#[cfg(target_os = "macos")]
fn normalize_launch_target(path: PathBuf) -> Option<LaunchTarget> {
if path
.extension()
.and_then(|value| value.to_str())
.is_some_and(|value| value.eq_ignore_ascii_case("app"))
{
let exec_path = path.join("Contents").join("MacOS").join(MACOS_PET_APP_NAME);
return exec_path.exists().then_some(LaunchTarget { exec_path });
}
path.exists().then_some(LaunchTarget { exec_path: path })
}
#[cfg(target_os = "windows")]
fn normalize_launch_target(path: PathBuf) -> Option<LaunchTarget> {
path.exists().then_some(LaunchTarget { exec_path: path })
}
#[cfg(not(any(target_os = "macos", target_os = "windows")))]
fn normalize_launch_target(path: PathBuf) -> Option<LaunchTarget> {
path.exists().then_some(LaunchTarget { exec_path: path })
}
async fn companion_pet_ws(
ws: WebSocketUpgrade,
State(router_state): State<CompanionRouterState>,
) -> impl IntoResponse {
ws.on_upgrade(move |socket| handle_companion_socket(router_state, socket))
}
async fn handle_companion_socket(router_state: CompanionRouterState, socket: WebSocket) {
let connection_id = Uuid::new_v4().to_string();
let (mut writer, mut reader) = socket.split();
let (tx, mut rx) = mpsc::unbounded_channel::<String>();
router_state
.service
.attach_sender(connection_id.clone(), tx)
.await;
router_state
.service
.mark_connection_open(connection_id.clone())
.await;
let writer_service = router_state.service.clone();
let writer_connection_id = connection_id.clone();
let writer_task = tokio::spawn(async move {
while let Some(message) = rx.recv().await {
if writer.send(Message::Text(message.into())).await.is_err() {
break;
}
}
writer_service
.mark_connection_closed(&writer_connection_id, "桌宠连接写入通道已关闭".to_string())
.await;
});
while let Some(message) = reader.next().await {
match message {
Ok(Message::Text(text)) => {
router_state
.service
.handle_incoming_message(&router_state.app_handle, &connection_id, &text)
.await;
}
Ok(Message::Binary(_)) => {}
Ok(Message::Ping(_)) => {}
Ok(Message::Pong(_)) => {}
Ok(Message::Close(_)) => {
break;
}
Err(error) => {
router_state
.service
.mark_connection_closed(&connection_id, format!("桌宠连接读取失败: {error}"))
.await;
writer_task.abort();
return;
}
}
}
writer_task.abort();
router_state
.service
.mark_connection_closed(&connection_id, "桌宠连接已关闭".to_string())
.await;
}
pub async fn get_pet_status_global(
state: &CompanionServiceState,
) -> Result<CompanionPetStatus, String> {
Ok(state.snapshot().await)
}
pub async fn launch_pet_global(
state: &CompanionServiceState,
request: CompanionLaunchPetRequest,
) -> Result<CompanionLaunchPetResult, String> {
state.launch_pet(request).await
}
pub async fn send_pet_command_global(
state: &CompanionServiceState,
request: CompanionPetCommandRequest,
) -> Result<CompanionPetSendResult, String> {
state.send_pet_command(request).await
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use tokio::sync::mpsc;
#[tokio::test]
async fn send_pet_command_without_active_sender_returns_not_delivered() {
let state = CompanionServiceState::default();
let result = state
.send_pet_command(CompanionPetCommandRequest {
event: "pet.provider_overview".to_string(),
payload: json!({
"total_provider_count": 2
}),
})
.await
.unwrap();
assert!(!result.delivered);
assert!(!result.connected);
}
#[tokio::test]
async fn send_pet_command_serializes_envelope_and_updates_visual_state() {
let state = CompanionServiceState::default();
let (tx, mut rx) = mpsc::unbounded_channel::<String>();
let connection_id = "conn-1".to_string();
state.attach_sender(connection_id.clone(), tx).await;
state.mark_connection_open(connection_id).await;
let result = state
.send_pet_command(CompanionPetCommandRequest {
event: "pet.state_changed".to_string(),
payload: json!({
"state": "thinking",
"total_provider_count": 3
}),
})
.await
.unwrap();
assert!(result.delivered);
assert!(result.connected);
let outbound = rx.recv().await.expect("应收到发送给桌宠的消息");
let envelope: serde_json::Value = serde_json::from_str(&outbound).expect("应输出合法 JSON");
assert_eq!(
envelope["protocol_version"],
json!(DEFAULT_PROTOCOL_VERSION)
);
assert_eq!(envelope["event"], json!("pet.state_changed"));
assert_eq!(envelope["payload"]["state"], json!("thinking"));
assert_eq!(envelope["payload"]["total_provider_count"], json!(3));
let snapshot = state.snapshot().await;
assert_eq!(snapshot.last_event.as_deref(), Some("pet.state_changed"));
assert_eq!(snapshot.last_state, Some(CompanionPetVisualState::Thinking));
assert!(snapshot.connected);
}
}
+1
View File
@@ -20,6 +20,7 @@ pub mod browser_profile_service;
pub mod browser_runtime_window;
pub mod chat_history_service;
pub mod claw_solution_service;
pub mod companion_service;
pub mod conversation_statistics_service;
pub mod environment_service;
pub mod execution_tracker_service;
+1
View File
@@ -42,6 +42,7 @@ const SITE_SEARCH_SKILL_CONTENT: &str =
const SITE_SEARCH_ADAPTER_CATALOG_CONTENT: &str =
include_str!("../../resources/default-skills/site_search/references/adapter-catalog.md");
#[cfg(test)]
const BUNDLED_SITE_ADAPTER_INDEX_CONTENT: &str =
include_str!("../../resources/site-adapters/bundled/index.json");