From 9d7c1f2e52a23a43e8e6560cd2f99ce9ee24bb01 Mon Sep 17 00:00:00 2001 From: oudecheng <13802883547@139.com> Date: Mon, 6 Jul 2026 18:07:33 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E9=9B=86=E6=88=90=E4=B8=93=E5=AE=B6?= =?UTF-8?q?=E6=8F=90=E7=A4=BA=E8=AF=8D=E5=88=B0=E7=BD=91=E5=85=B3=E9=85=8D?= =?UTF-8?q?=E7=BD=AE=E3=80=81=E4=BE=9D=E8=B5=96=E6=B3=A8=E5=85=A5=E4=B8=8E?= =?UTF-8?q?HTTP=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - config: 新增 ExpertsConfig,默认 sources=[user,project], enabled=true - AgentFactory: 新增 experts 字段,providers 顺序为 Agent→Skill→Expert→Subagent→Todo - GatewayState: 持有 Arc,from_config 中初始化 - runtime: build_session_manager_with_sender 接收 experts 参数 - 路由: 注册 7 个专家 API (list/toggle/create/update/delete/selected/select) - http: 实现 7 个处理器,支持专家 CRUD 与 session 级选择 - session/cli: 测试便利构造器与 empty_config/build_config 补齐 experts 字段 --- src/cli/init.rs | 2 + src/config/mod.rs | 30 ++++ src/gateway/agent_factory.rs | 7 + src/gateway/http.rs | 340 ++++++++++++++++++++++++++++++++++- src/gateway/mod.rs | 18 ++ src/gateway/runtime.rs | 4 + src/gateway/session.rs | 8 + 7 files changed, 408 insertions(+), 1 deletion(-) diff --git a/src/cli/init.rs b/src/cli/init.rs index 3337fd3..9c51659 100644 --- a/src/cli/init.rs +++ b/src/cli/init.rs @@ -79,6 +79,7 @@ impl InitWizard { mcp_servers: HashMap::new(), image_context: crate::config::ImageContextConfig::default(), subagents: crate::config::SubagentsConfig::default(), + experts: crate::config::ExpertsConfig::default(), } } @@ -832,6 +833,7 @@ impl InitWizard { mcp_servers: existing.mcp_servers.clone(), image_context: existing.image_context.clone(), subagents: existing.subagents.clone(), + experts: existing.experts.clone(), } } diff --git a/src/config/mod.rs b/src/config/mod.rs index 23ecc1a..77b655d 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -38,6 +38,8 @@ pub struct Config { pub image_context: ImageContextConfig, #[serde(default)] pub subagents: SubagentsConfig, + #[serde(default)] + pub experts: ExpertsConfig, } /// 图片上下文限制配置 @@ -169,6 +171,34 @@ impl Default for SubagentsConfig { } } +/// 专家提示词配置 +#[derive(Debug, Clone, Deserialize, Serialize)] +pub struct ExpertsConfig { + /// 是否启用专家发现与注入 + #[serde(default = "default_experts_enabled")] + pub enabled: bool, + /// 定义来源优先级 + #[serde(default = "default_experts_sources")] + pub sources: Vec, +} + +fn default_experts_enabled() -> bool { + true +} + +fn default_experts_sources() -> Vec { + vec!["user".to_string(), "project".to_string()] +} + +impl Default for ExpertsConfig { + fn default() -> Self { + Self { + enabled: default_experts_enabled(), + sources: default_experts_sources(), + } + } +} + #[derive(Debug, Clone, Default, Deserialize, Serialize)] pub struct ToolsConfig { #[serde(default)] diff --git a/src/gateway/agent_factory.rs b/src/gateway/agent_factory.rs index 0854f8f..9c508b0 100644 --- a/src/gateway/agent_factory.rs +++ b/src/gateway/agent_factory.rs @@ -2,6 +2,8 @@ use std::sync::Arc; use crate::agent::{AgentError, AgentLoop, CompositeSystemPromptProvider}; use crate::config::LLMProviderConfig; +use crate::experts::ExpertPromptProvider; +use crate::experts::ExpertRuntime; use crate::gateway::agent_prompt_provider::AgentPromptProvider; use crate::gateway::todo_prompt_provider::TodoPromptProvider; use crate::skills::{SkillPromptProvider, SkillRuntime}; @@ -14,6 +16,7 @@ use crate::tools::{ToolContext, ToolRegistry}; pub(crate) struct AgentFactory { tools: Arc, skills: Arc, + experts: Arc, subagent_runtime: Arc, reinject_every: usize, prompt_repository: Arc, @@ -38,6 +41,7 @@ impl AgentFactory { pub(crate) fn new( tools: Arc, skills: Arc, + experts: Arc, subagent_runtime: Arc, reinject_every: usize, prompt_repository: Arc, @@ -52,6 +56,7 @@ impl AgentFactory { Self { tools, skills, + experts, subagent_runtime, reinject_every, prompt_repository, @@ -74,6 +79,7 @@ impl AgentFactory { ); // 创建组合的系统提示词提供者 + // 顺序:AgentPrompt → SkillPrompt → ExpertPrompt → SubagentPrompt → TodoPrompt let system_prompt_provider = Arc::new(CompositeSystemPromptProvider::new(vec![ Box::new(AgentPromptProvider::new( self.reinject_every, @@ -81,6 +87,7 @@ impl AgentFactory { self.prompt_repository.clone(), )), Box::new(SkillPromptProvider::new(self.skills.clone())), + Box::new(ExpertPromptProvider::new(self.experts.clone())), Box::new(SubagentPromptProvider::new(self.subagent_runtime.clone())), Box::new(TodoPromptProvider::new()), ])); diff --git a/src/gateway/http.rs b/src/gateway/http.rs index f050890..eccbf88 100644 --- a/src/gateway/http.rs +++ b/src/gateway/http.rs @@ -1,10 +1,11 @@ -use axum::{Json, extract::State}; +use axum::{Json, extract::{Query, State}}; use axum::http::StatusCode; use serde::{Deserialize, Serialize}; use std::sync::Arc; use super::GatewayState; use crate::config::{Config, get_default_config_path}; +use crate::experts::{Expert, ExpertScope, ExpertWithStatus}; use crate::skills::SkillWithStatus; use crate::tools::task::runtime::{SubagentScope, SubagentWithStatus}; @@ -410,3 +411,340 @@ pub async fn subagents_toggle( } } } + +// ===================== Experts ===================== + +#[derive(Deserialize)] +pub struct ExpertToggleRequest { + pub name: String, + pub scope: String, + pub enabled: bool, +} + +#[derive(Serialize)] +pub struct ExpertToggleResponse { + success: bool, + #[serde(skip_serializing_if = "Option::is_none")] + changed: Option, + #[serde(skip_serializing_if = "Option::is_none")] + available: Option, + #[serde(skip_serializing_if = "Option::is_none")] + disabled_in_scopes: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + error: Option, +} + +#[derive(Serialize)] +pub struct ExpertListResponse { + experts_system_enabled: bool, + total: usize, + experts: Vec, +} + +#[derive(Deserialize)] +pub struct ExpertCreateRequest { + pub name: String, + pub description: String, + pub body: String, + pub scope: String, +} + +#[derive(Deserialize)] +pub struct ExpertUpdateRequest { + pub name: String, + pub scope: String, + pub description: Option, + pub body: Option, +} + +#[derive(Deserialize)] +pub struct ExpertDeleteRequest { + pub name: String, + pub scope: String, +} + +#[derive(Deserialize)] +pub struct ExpertSelectedQuery { + pub session_id: String, +} + +#[derive(Serialize)] +pub struct ExpertSelectedResponse { + pub expert_name: Option, + pub expert: Option, +} + +#[derive(Deserialize)] +pub struct ExpertSelectRequest { + pub session_id: String, + pub expert_name: Option, +} + +#[derive(Serialize)] +pub struct ExpertSelectResponse { + pub success: bool, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +#[derive(Serialize)] +pub struct ExpertResponse { + pub name: String, + pub description: String, + pub body: String, + pub source: String, + pub path: String, +} + +impl From for ExpertResponse { + fn from(expert: Expert) -> Self { + Self { + name: expert.name, + description: expert.description, + body: expert.body, + source: expert.source.as_str().to_string(), + path: expert.path.display().to_string(), + } + } +} + +#[derive(Serialize)] +pub struct ExpertDeleteResponse { + pub success: bool, + pub path: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// GET /api/experts — Return all discovered experts with their disabled status +pub async fn experts_list( + State(state): State>, +) -> Json { + let experts_enabled = state.config.read().await.experts.enabled; + + if !experts_enabled { + return Json(ExpertListResponse { + experts_system_enabled: false, + total: 0, + experts: vec![], + }); + } + + let experts = state.experts.list_experts_with_status(); + let total = experts.len(); + + Json(ExpertListResponse { + experts_system_enabled: true, + total, + experts, + }) +} + +/// POST /api/experts/toggle — Enable or disable a specific expert +pub async fn experts_toggle( + State(state): State>, + Json(req): Json, +) -> (StatusCode, Json) { + let scope = match ExpertScope::parse(&req.scope) { + Some(s) => s, + None => { + return ( + StatusCode::BAD_REQUEST, + Json(ExpertToggleResponse { + success: false, + changed: None, + available: None, + disabled_in_scopes: None, + error: Some(format!("invalid scope: {}", req.scope)), + }), + ); + } + }; + + let result = if req.enabled { + state.experts.enable_expert(scope, &req.name) + } else { + state.experts.disable_expert(scope, &req.name) + }; + + match result { + Ok(change) => ( + StatusCode::OK, + Json(ExpertToggleResponse { + success: true, + changed: Some(change.changed), + available: Some(change.available), + disabled_in_scopes: Some( + change + .disabled_in_scopes + .iter() + .map(|s| s.as_str().to_string()) + .collect(), + ), + error: None, + }), + ), + Err(msg) => { + let status = if msg.contains("not found") { + StatusCode::NOT_FOUND + } else { + StatusCode::INTERNAL_SERVER_ERROR + }; + ( + status, + Json(ExpertToggleResponse { + success: false, + changed: None, + available: None, + disabled_in_scopes: None, + error: Some(msg), + }), + ) + } + } +} + +/// POST /api/experts/create — Create a new expert +pub async fn experts_create( + State(state): State>, + Json(req): Json, +) -> Result, (StatusCode, String)> { + let scope = ExpertScope::parse(&req.scope) + .ok_or_else(|| (StatusCode::BAD_REQUEST, format!("invalid scope: {}", req.scope)))?; + + let expert = state + .experts + .create_expert(scope, &req.name, &req.description, &req.body, true) + .map_err(|err| { + let status = if err.contains("already exists") { + StatusCode::CONFLICT + } else { + StatusCode::BAD_REQUEST + }; + (status, err) + })?; + + Ok(Json(ExpertResponse::from(expert))) +} + +/// PUT /api/experts/update — Update an existing expert +pub async fn experts_update( + State(state): State>, + Json(req): Json, +) -> Result, (StatusCode, String)> { + let scope = ExpertScope::parse(&req.scope) + .ok_or_else(|| (StatusCode::BAD_REQUEST, format!("invalid scope: {}", req.scope)))?; + + let expert = state + .experts + .update_expert( + scope, + &req.name, + req.description.as_deref(), + req.body.as_deref(), + true, + ) + .map_err(|err| { + let status = if err.contains("not found") { + StatusCode::NOT_FOUND + } else { + StatusCode::BAD_REQUEST + }; + (status, err) + })?; + + Ok(Json(ExpertResponse::from(expert))) +} + +/// DELETE /api/experts/delete?name=&scope= — Delete an expert +pub async fn experts_delete( + State(state): State>, + Query(req): Query, +) -> Result, (StatusCode, Json)> { + let scope = match ExpertScope::parse(&req.scope) { + Some(s) => s, + None => { + return Err(( + StatusCode::BAD_REQUEST, + Json(ExpertDeleteResponse { + success: false, + path: String::new(), + error: Some(format!("invalid scope: {}", req.scope)), + }), + )); + } + }; + + match state.experts.delete_expert(scope, &req.name, true) { + Ok(path) => Ok(Json(ExpertDeleteResponse { + success: true, + path: path.display().to_string(), + error: None, + })), + Err(msg) => { + let status = if msg.contains("not found") { + StatusCode::NOT_FOUND + } else { + StatusCode::INTERNAL_SERVER_ERROR + }; + Err(( + status, + Json(ExpertDeleteResponse { + success: false, + path: String::new(), + error: Some(msg), + }), + )) + } + } +} + +/// GET /api/experts/selected?session_id=... — Return the currently selected expert for a session +pub async fn experts_selected( + State(state): State>, + Query(q): Query, +) -> Json { + let expert_name = state.experts.selected_expert_name_for(&q.session_id); + + let expert = expert_name.as_ref().and_then(|name| { + // Build an ExpertWithStatus from the discovered catalog. + state + .experts + .list_experts_with_status() + .into_iter() + .find(|e| &e.name == name) + }); + + Json(ExpertSelectedResponse { + expert_name, + expert, + }) +} + +/// POST /api/experts/select — Select (or clear) the expert for a session +pub async fn experts_select( + State(state): State>, + Json(req): Json, +) -> (StatusCode, Json) { + let result = match req.expert_name { + Some(name) => state.experts.select_expert(&req.session_id, &name), + None => state.experts.clear_expert(&req.session_id), + }; + + match result { + Ok(()) => ( + StatusCode::OK, + Json(ExpertSelectResponse { + success: true, + error: None, + }), + ), + Err(msg) => ( + StatusCode::BAD_REQUEST, + Json(ExpertSelectResponse { + success: false, + error: Some(msg), + }), + ), + } +} diff --git a/src/gateway/mod.rs b/src/gateway/mod.rs index 886f868..f64e66e 100644 --- a/src/gateway/mod.rs +++ b/src/gateway/mod.rs @@ -65,6 +65,7 @@ pub struct GatewayState { pub restart_tx: watch::Sender, pub mcp_manager: Option>, pub skills: Arc, + pub experts: Arc, pub subagent_runtime: Arc, } @@ -83,6 +84,7 @@ impl GatewayState { let session_ttl_hours = config.gateway.session_ttl_hours; let skills = Arc::new(SkillRuntime::from_config(config.skills.clone())); + let experts = Arc::new(crate::experts::ExpertRuntime::from_config(config.experts.clone())); let channel_manager = ChannelManager::new(); let bus = channel_manager.bus(); @@ -97,6 +99,7 @@ impl GatewayState { provider_config, provider_configs, skills.clone(), + experts.clone(), Arc::new(BusSessionMessageSender::new(bus.clone())), std::collections::HashSet::new(), config.tools.task.clone(), @@ -125,6 +128,7 @@ impl GatewayState { restart_tx, mcp_manager, skills, + experts, subagent_runtime, }) } @@ -239,6 +243,13 @@ pub async fn run( .route("/api/skills/toggle", routing::post(http::skills_toggle)) .route("/api/subagents", routing::get(http::subagents_list)) .route("/api/subagents/toggle", routing::post(http::subagents_toggle)) + .route("/api/experts", routing::get(http::experts_list)) + .route("/api/experts/toggle", routing::post(http::experts_toggle)) + .route("/api/experts/create", routing::post(http::experts_create)) + .route("/api/experts/update", routing::put(http::experts_update)) + .route("/api/experts/delete", routing::delete(http::experts_delete)) + .route("/api/experts/selected", routing::get(http::experts_selected)) + .route("/api/experts/select", routing::post(http::experts_select)) .route("/ws", routing::get(ws::ws_handler)) .fallback(static_handler) .with_state(state.clone()) @@ -253,6 +264,13 @@ pub async fn run( .route("/api/skills/toggle", routing::post(http::skills_toggle)) .route("/api/subagents", routing::get(http::subagents_list)) .route("/api/subagents/toggle", routing::post(http::subagents_toggle)) + .route("/api/experts", routing::get(http::experts_list)) + .route("/api/experts/toggle", routing::post(http::experts_toggle)) + .route("/api/experts/create", routing::post(http::experts_create)) + .route("/api/experts/update", routing::put(http::experts_update)) + .route("/api/experts/delete", routing::delete(http::experts_delete)) + .route("/api/experts/selected", routing::get(http::experts_selected)) + .route("/api/experts/select", routing::post(http::experts_select)) .route("/ws", routing::get(ws::ws_handler)) .fallback_service(ServeDir::new(&static_dir)) .with_state(state.clone()) diff --git a/src/gateway/runtime.rs b/src/gateway/runtime.rs index 6f7296b..7d7ec76 100644 --- a/src/gateway/runtime.rs +++ b/src/gateway/runtime.rs @@ -45,6 +45,7 @@ pub(crate) fn build_session_manager( provider_config: LLMProviderConfig, provider_configs: HashMap, skills: Arc, + experts: Arc, disabled_tools: HashSet, task_config: TaskConfig, subagents_config: SubagentsConfig, @@ -60,6 +61,7 @@ pub(crate) fn build_session_manager( provider_config, provider_configs, skills, + experts, Arc::new(NoopSessionMessageSender), disabled_tools, task_config, @@ -79,6 +81,7 @@ pub(crate) fn build_session_manager_with_sender( provider_config: LLMProviderConfig, provider_configs: HashMap, skills: Arc, + experts: Arc, session_message_sender: Arc, disabled_tools: HashSet, task_config: TaskConfig, @@ -269,6 +272,7 @@ pub(crate) fn build_session_manager_with_sender( let agent_factory = AgentFactory::new( tools.clone(), skills.clone(), + experts.clone(), subagent_runtime.clone(), agent_prompt_reinject_every as usize, prompt_repository.clone(), diff --git a/src/gateway/session.rs b/src/gateway/session.rs index fcfef3a..cdc415b 100644 --- a/src/gateway/session.rs +++ b/src/gateway/session.rs @@ -255,9 +255,13 @@ impl Session { let conversations: Arc = store.clone(); let skill_events: Arc = store.clone(); let prompt_repository: Arc = store.clone(); + let experts = Arc::new(crate::experts::ExpertRuntime::from_config( + crate::config::ExpertsConfig::default(), + )); let agent_factory = AgentFactory::new( tools, skills.clone(), + experts, subagent_runtime, agent_prompt_reinject_every as usize, prompt_repository.clone(), @@ -670,6 +674,9 @@ impl SessionManager { session_ttl_hours: Option, mcp_config: crate::mcp::McpConfig, ) -> Result { + let experts = Arc::new(crate::experts::ExpertRuntime::from_config( + crate::config::ExpertsConfig::default(), + )); super::runtime::build_session_manager( agent_prompt_reinject_every, show_tool_results, @@ -677,6 +684,7 @@ impl SessionManager { provider_config, provider_configs, skills, + experts, disabled_tools, task_config, subagents_config,