refactor: 迁移 MCP 模块到 proxycast-mcp crate

- 创建 proxycast-mcp crate(client, manager, tool_converter, types)
- 主 crate mcp/mod.rs 替换为 re-export 层
- McpClientManager、ProxyCastMcpClient 等类型迁移
This commit is contained in:
coso
2026-02-08 22:41:07 +08:00
parent f169e93e6e
commit bb319d7bb0
7 changed files with 111 additions and 193 deletions
+17
View File
@@ -0,0 +1,17 @@
[package]
name = "proxycast-mcp"
version.workspace = true
edition.workspace = true
authors.workspace = true
repository.workspace = true
[dependencies]
proxycast-core.workspace = true
serde.workspace = true
serde_json.workspace = true
tokio.workspace = true
async-trait.workspace = true
tracing.workspace = true
thiserror.workspace = true
glob.workspace = true
rmcp.workspace = true
@@ -1,13 +1,11 @@
//! MCP 客户端实现
//!
//! 本模块实现 rmcp 的 ClientHandler trait,处理:
//! - 客户端信息返回
//! - 进度通知处理
//! - 日志消息处理
//! - 与 Tauri 事件系统的集成
//! 实现 rmcp 的 ClientHandler trait,处理通知和回调。
//! 使用 DynEmitter 替代 Tauri AppHandle 进行事件发射。
#![allow(dead_code)]
use proxycast_core::DynEmitter;
use rmcp::{
model::{
ClientCapabilities, ClientInfo, Implementation, LoggingMessageNotification,
@@ -18,7 +16,6 @@ use rmcp::{
ClientHandler, RoleClient,
};
use std::sync::Arc;
use tauri::Emitter;
use tokio::sync::{mpsc, Mutex};
use tracing::{debug, info, warn};
@@ -42,67 +39,49 @@ pub struct McpLogMessagePayload {
}
/// ProxyCast MCP 客户端处理器
///
/// 实现 rmcp::ClientHandler trait,处理 MCP 服务器的通知和回调
pub struct ProxyCastMcpClient {
/// Tauri AppHandle(用于发送事件)
app_handle: Option<tauri::AppHandle>,
/// 服务器名称(用于事件标识)
emitter: Option<DynEmitter>,
server_name: String,
/// 通知订阅者(用于内部通知分发)
notification_handlers: Arc<Mutex<Vec<mpsc::Sender<ServerNotification>>>>,
}
impl ProxyCastMcpClient {
/// 创建新的 MCP 客户端处理器
///
/// # Arguments
/// * `server_name` - MCP 服务器名称,用于事件标识
/// * `app_handle` - Tauri AppHandle,用于发送事件到前端
pub fn new(server_name: String, app_handle: Option<tauri::AppHandle>) -> Self {
pub fn new(server_name: String, emitter: Option<DynEmitter>) -> Self {
Self {
app_handle,
emitter,
server_name,
notification_handlers: Arc::new(Mutex::new(Vec::new())),
}
}
/// 获取通知处理器的引用(用于订阅通知)
pub fn notification_handlers(&self) -> Arc<Mutex<Vec<mpsc::Sender<ServerNotification>>>> {
self.notification_handlers.clone()
}
/// 订阅服务器通知
///
/// 返回一个接收器,用于接收来自 MCP 服务器的通知
pub async fn subscribe(&self) -> mpsc::Receiver<ServerNotification> {
let (tx, rx) = mpsc::channel(16);
self.notification_handlers.lock().await.push(tx);
rx
}
/// 发送 Tauri 事件到前端
fn emit_event<T: serde::Serialize + Clone>(&self, event: &str, payload: T) {
if let Some(ref app_handle) = self.app_handle {
if let Err(e) = app_handle.emit(event, payload) {
warn!(
server_name = %self.server_name,
event = %event,
error = %e,
"发送 Tauri 事件失败"
);
/// 发送事件(通过 DynEmitter)
fn emit_event<T: serde::Serialize>(&self, event: &str, payload: &T) {
if let Some(ref emitter) = self.emitter {
if let Ok(value) = serde_json::to_value(payload) {
if let Err(e) = emitter.emit_event(event, &value) {
warn!(
server_name = %self.server_name,
event = %event,
error = %e,
"发送事件失败"
);
}
}
}
}
}
impl ClientHandler for ProxyCastMcpClient {
/// 返回客户端信息
///
/// 提供 ProxyCast 客户端的标识信息,包括:
/// - 协议版本
/// - 客户端能力(采样支持)
/// - 客户端实现信息
fn get_info(&self) -> ClientInfo {
ClientInfo {
protocol_version: ProtocolVersion::V_2025_03_26,
@@ -117,19 +96,11 @@ impl ClientHandler for ProxyCastMcpClient {
}
}
/// 处理进度通知
///
/// 当 MCP 服务器发送进度更新时调用此方法。
/// 进度信息会:
/// 1. 记录到日志
/// 2. 发送到 Tauri 事件系统(前端可监听)
/// 3. 分发给内部通知订阅者
async fn on_progress(
&self,
params: ProgressNotificationParam,
context: NotificationContext<RoleClient>,
) {
// 记录进度日志
debug!(
server_name = %self.server_name,
progress_token = ?params.progress_token,
@@ -138,7 +109,6 @@ impl ClientHandler for ProxyCastMcpClient {
"收到 MCP 进度通知"
);
// 发送 Tauri 事件到前端
let payload = McpProgressPayload {
server_name: self.server_name.clone(),
progress_token: format!("{:?}", params.progress_token),
@@ -146,9 +116,8 @@ impl ClientHandler for ProxyCastMcpClient {
total: params.total,
message: None,
};
self.emit_event("mcp:progress", payload);
self.emit_event("mcp:progress", &payload);
// 分发给内部通知订阅者
let notification = ServerNotification::ProgressNotification(ProgressNotification {
params: params.clone(),
method: ProgressNotificationMethod,
@@ -161,97 +130,38 @@ impl ClientHandler for ProxyCastMcpClient {
}
}
/// 处理日志消息通知
///
/// 当 MCP 服务器发送日志消息时调用此方法。
/// 日志消息会:
/// 1. 根据级别记录到本地日志
/// 2. 发送到 Tauri 事件系统(前端可监听)
/// 3. 分发给内部通知订阅者
async fn on_logging_message(
&self,
params: LoggingMessageNotificationParam,
context: NotificationContext<RoleClient>,
) {
// 根据日志级别记录
let level_str = format!("{:?}", params.level);
match params.level {
rmcp::model::LoggingLevel::Debug => {
debug!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [DEBUG]"
);
debug!(server_name = %self.server_name, logger = ?params.logger, data = ?params.data, "MCP 服务器日志 [DEBUG]");
}
rmcp::model::LoggingLevel::Info => {
info!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [INFO]"
);
info!(server_name = %self.server_name, logger = ?params.logger, data = ?params.data, "MCP 服务器日志 [INFO]");
}
rmcp::model::LoggingLevel::Notice => {
info!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [NOTICE]"
);
info!(server_name = %self.server_name, logger = ?params.logger, data = ?params.data, "MCP 服务器日志 [NOTICE]");
}
rmcp::model::LoggingLevel::Warning => {
warn!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [WARNING]"
);
warn!(server_name = %self.server_name, logger = ?params.logger, data = ?params.data, "MCP 服务器日志 [WARNING]");
}
rmcp::model::LoggingLevel::Error => {
tracing::error!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [ERROR]"
);
}
rmcp::model::LoggingLevel::Critical => {
tracing::error!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [CRITICAL]"
);
}
rmcp::model::LoggingLevel::Alert => {
tracing::error!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [ALERT]"
);
}
rmcp::model::LoggingLevel::Emergency => {
tracing::error!(
server_name = %self.server_name,
logger = ?params.logger,
data = ?params.data,
"MCP 服务器日志 [EMERGENCY]"
);
_ => {
tracing::error!(server_name = %self.server_name, logger = ?params.logger, data = ?params.data, level = %level_str, "MCP 服务器日志");
}
}
// 发送 Tauri 事件到前端
let payload = McpLogMessagePayload {
server_name: self.server_name.clone(),
level: level_str,
logger: params.logger.clone(),
data: params.data.clone(),
};
self.emit_event("mcp:log_message", payload);
self.emit_event("mcp:log_message", &payload);
// 分发给内部通知订阅者
let notification =
ServerNotification::LoggingMessageNotification(LoggingMessageNotification {
params: params.clone(),
@@ -267,32 +177,23 @@ impl ClientHandler for ProxyCastMcpClient {
}
/// MCP 客户端包装器
///
/// 封装 rmcp 客户端和相关状态
pub struct McpClientWrapper {
/// 服务器名称
pub server_name: String,
/// 服务器配置
pub config: super::types::McpServerConfig,
/// 子进程句柄
pub process: Option<tokio::process::Child>,
/// 服务器能力信息
pub server_info: Option<super::types::McpServerCapabilities>,
/// 客户端处理器
pub client_handler: Arc<ProxyCastMcpClient>,
/// rmcp 运行服务(用于发送请求)
pub running_service:
Option<rmcp::service::RunningService<rmcp::RoleClient, ProxyCastMcpClient>>,
}
impl McpClientWrapper {
/// 创建新的客户端包装器
pub fn new(
server_name: String,
config: super::types::McpServerConfig,
app_handle: Option<tauri::AppHandle>,
emitter: Option<DynEmitter>,
) -> Self {
let client_handler = Arc::new(ProxyCastMcpClient::new(server_name.clone(), app_handle));
let client_handler = Arc::new(ProxyCastMcpClient::new(server_name.clone(), emitter));
Self {
server_name,
@@ -304,22 +205,18 @@ impl McpClientWrapper {
}
}
/// 获取客户端处理器的引用
pub fn handler(&self) -> Arc<ProxyCastMcpClient> {
self.client_handler.clone()
}
/// 设置子进程句柄
pub fn set_process(&mut self, process: tokio::process::Child) {
self.process = Some(process);
}
/// 设置服务器能力信息
pub fn set_server_info(&mut self, info: super::types::McpServerCapabilities) {
self.server_info = Some(info);
}
/// 设置 rmcp 运行服务
pub fn set_running_service(
&mut self,
service: rmcp::service::RunningService<rmcp::RoleClient, ProxyCastMcpClient>,
@@ -327,14 +224,12 @@ impl McpClientWrapper {
self.running_service = Some(service);
}
/// 获取 rmcp 运行服务的引用
pub fn running_service(
&self,
) -> Option<&rmcp::service::RunningService<rmcp::RoleClient, ProxyCastMcpClient>> {
self.running_service.as_ref()
}
/// 终止子进程
pub async fn kill_process(&mut self) -> Result<(), std::io::Error> {
if let Some(ref mut process) = self.process {
process.kill().await?;
@@ -355,7 +250,6 @@ mod tests {
let info = client.get_info();
assert_eq!(info.client_info.name, "proxycast");
assert_eq!(info.client_info.version, env!("CARGO_PKG_VERSION"));
assert_eq!(
info.client_info.title,
Some("ProxyCast MCP Client".to_string())
@@ -373,7 +267,7 @@ mod tests {
timeout: 30,
};
let wrapper = McpClientWrapper::new("test-server".to_string(), config.clone(), None);
let wrapper = McpClientWrapper::new("test-server".to_string(), config, None);
assert_eq!(wrapper.server_name, "test-server");
assert_eq!(wrapper.config.command, "test-command");
@@ -385,15 +279,12 @@ mod tests {
async fn test_notification_subscription() {
let client = ProxyCastMcpClient::new("test-server".to_string(), None);
// 订阅通知
let mut rx = client.subscribe().await;
// 验证订阅者已添加
let handlers = client.notification_handlers.lock().await;
assert_eq!(handlers.len(), 1);
drop(handlers);
// 验证接收器可用(不会阻塞)
assert!(rx.try_recv().is_err()); // 应该是空的
assert!(rx.try_recv().is_err());
}
}
+20
View File
@@ -0,0 +1,20 @@
//! ProxyCast MCP Crate
//!
//! MCP(Model Context Protocol)集成模块,提供 MCP 协议的客户端实现。
//! 使用 DynEmitter 替代 Tauri AppHandle 进行事件发射,实现与 Tauri 的解耦。
pub mod client;
pub mod manager;
pub mod tool_converter;
pub mod types;
pub use client::{McpClientWrapper, ProxyCastMcpClient};
pub use manager::McpClientManager;
pub use tool_converter::ToolConverter;
pub use types::{
McpContent, McpError, McpManagerState, McpPromptArgument, McpPromptDefinition,
McpPromptMessage, McpPromptResult, McpResourceContent, McpResourceDefinition,
McpServerCapabilities, McpServerConfig, McpServerErrorPayload, McpServerInfo,
McpServerStartedPayload, McpServerStoppedPayload, McpToolCall, McpToolDefinition,
McpToolResult, McpToolsUpdatedPayload,
};
@@ -26,11 +26,11 @@
#![allow(dead_code)]
use proxycast_core::DynEmitter;
use std::collections::HashMap;
use std::process::Stdio;
use std::sync::Arc;
use std::time::Duration;
use tauri::Emitter;
use tokio::io::AsyncReadExt;
use tokio::process::Command;
use tokio::sync::RwLock;
@@ -39,8 +39,8 @@ use tracing::{debug, error, info, warn};
use rmcp::transport::TokioChildProcess;
use rmcp::ServiceExt;
use super::client::McpClientWrapper;
use super::types::*;
use crate::client::McpClientWrapper;
use crate::types::*;
/// MCP 客户端管理器
///
@@ -89,14 +89,14 @@ pub struct McpClientManager {
/// - Some(tools): 缓存有效
tool_cache: Arc<RwLock<Option<Vec<McpToolDefinition>>>>,
/// Tauri AppHandle(用于发送事件)
/// 事件发射器
///
/// 用于向前端发送 MCP 相关事件,如:
/// - mcp:server_started
/// - mcp:server_stopped
/// - mcp:server_error
/// - mcp:tools_updated
app_handle: Option<tauri::AppHandle>,
emitter: Option<DynEmitter>,
}
impl McpClientManager {
@@ -104,24 +104,24 @@ impl McpClientManager {
///
/// # Arguments
///
/// * `app_handle` - Tauri AppHandle,用于发送事件到前端。
/// * `emitter` - 事件发射器,用于发送事件到前端。
/// 如果为 None,则不会发送事件。
///
/// # Returns
///
/// 返回初始化的 McpClientManager 实例,连接池和缓存均为空。
pub fn new(app_handle: Option<tauri::AppHandle>) -> Self {
pub fn new(emitter: Option<DynEmitter>) -> Self {
info!("创建 MCP 客户端管理器");
Self {
clients: Arc::new(RwLock::new(HashMap::new())),
tool_cache: Arc::new(RwLock::new(None)),
app_handle,
emitter,
}
}
/// 设置 Tauri AppHandle(用于发送前端事件)
pub fn set_app_handle(&mut self, app_handle: tauri::AppHandle) {
self.app_handle = Some(app_handle);
/// 设置事件发射器
pub fn set_emitter(&mut self, emitter: DynEmitter) {
self.emitter = Some(emitter);
}
// ========================================================================
@@ -292,15 +292,17 @@ impl McpClientManager {
/// * `event` - 事件名称
/// * `payload` - 事件数据
pub fn emit_event<T: serde::Serialize + Clone>(&self, event: &str, payload: T) {
if let Some(ref app_handle) = self.app_handle {
if let Err(e) = app_handle.emit(event, payload) {
warn!(
event = %event,
error = %e,
"发送 Tauri 事件失败"
);
} else {
debug!(event = %event, "发送 Tauri 事件");
if let Some(ref emitter) = self.emitter {
if let Ok(value) = serde_json::to_value(&payload) {
if let Err(e) = emitter.emit_event(event, &value) {
warn!(
event = %event,
error = %e,
"发送事件失败"
);
} else {
debug!(event = %event, "发送事件");
}
}
}
}
@@ -471,7 +473,7 @@ impl McpClientManager {
// 4. 初始化 MCP 客户端
let client_handler =
super::client::ProxyCastMcpClient::new(name.to_string(), self.app_handle.clone());
crate::client::ProxyCastMcpClient::new(name.to_string(), self.emitter.clone());
// 连接超时:至少 60 秒,避免 npx 首次下载时超时
let timeout_secs = std::cmp::max(config.timeout, 60);
@@ -538,10 +540,10 @@ impl McpClientManager {
});
// 创建客户端包装器
let mut wrapper = super::client::McpClientWrapper::new(
let mut wrapper = crate::client::McpClientWrapper::new(
name.to_string(),
config.clone(),
self.app_handle.clone(),
self.emitter.clone(),
);
if let Some(ref info) = server_info {
wrapper.set_server_info(info.clone());
@@ -1495,8 +1497,8 @@ impl McpClientManager {
pub type McpManagerState = Arc<tokio::sync::Mutex<McpClientManager>>;
/// 创建 MCP 管理器状态
pub fn create_mcp_manager_state(app_handle: Option<tauri::AppHandle>) -> McpManagerState {
Arc::new(tokio::sync::Mutex::new(McpClientManager::new(app_handle)))
pub fn create_mcp_manager_state(emitter: Option<DynEmitter>) -> McpManagerState {
Arc::new(tokio::sync::Mutex::new(McpClientManager::new(emitter)))
}
// ============================================================================
@@ -1528,7 +1530,7 @@ mod tests {
fn test_manager_creation() {
let manager = McpClientManager::new(None);
// 验证初始状态
assert!(manager.app_handle.is_none());
assert!(manager.emitter.is_none());
}
#[tokio::test]
@@ -232,13 +232,13 @@ pub struct McpToolsUpdatedPayload {
}
// ============================================================================
// Tauri 状态类型
// 状态类型
// ============================================================================
use std::sync::Arc;
use tokio::sync::Mutex;
/// MCP 客户端管理器状态(Tauri 托管状态)
/// MCP 客户端管理器状态
///
/// 使用 Arc<Mutex<McpClientManager>> 包装,支持跨线程共享和异步访问。
pub type McpManagerState = Arc<Mutex<super::manager::McpClientManager>>;
+14 -26
View File
@@ -1,31 +1,19 @@
//! MCP(Model Context Protocol)模块
//!
//! 本模块提供 MCP 协议的客户端实现,支持:
//! - MCP 服务器生命周期管理(启动、停止、状态监控)
//! - MCP 工具发现和调用
//! - MCP 提示词和资源访问
//! - 工具格式转换(OpenAI/Anthropic/Gemini)
//!
//! # 模块结构
//!
//! - `types`: MCP 数据类型定义
//! - `client`: MCP 客户端实现(rmcp ClientHandler)
//! - `manager`: MCP 客户端管理器(连接池、缓存)
//! - `tool_converter`: 工具格式转换器
//! 业务逻辑已迁移到 proxycast-mcp crate,
//! 本模块仅作为桥接层 re-export。
pub mod client;
pub mod manager;
pub mod tool_converter;
pub mod types;
// 从 proxycast-mcp crate re-export 所有公开类型
pub use proxycast_mcp::client;
pub use proxycast_mcp::manager;
pub use proxycast_mcp::tool_converter;
pub use proxycast_mcp::types;
// 显式导出,避免命名冲突
pub use client::{McpClientWrapper, ProxyCastMcpClient};
pub use manager::McpClientManager;
pub use tool_converter::ToolConverter;
pub use types::{
McpContent, McpError, McpManagerState, McpPromptArgument, McpPromptDefinition,
McpPromptMessage, McpPromptResult, McpResourceContent, McpResourceDefinition,
McpServerCapabilities, McpServerConfig, McpServerErrorPayload, McpServerInfo,
McpServerStartedPayload, McpServerStoppedPayload, McpToolCall, McpToolDefinition,
McpToolResult, McpToolsUpdatedPayload,
pub use proxycast_mcp::{McpClientManager, ProxyCastMcpClient};
pub use proxycast_mcp::{
McpClientWrapper, McpContent, McpError, McpManagerState, McpPromptArgument,
McpPromptDefinition, McpPromptMessage, McpPromptResult, McpResourceContent,
McpResourceDefinition, McpServerCapabilities, McpServerConfig, McpServerErrorPayload,
McpServerInfo, McpServerStartedPayload, McpServerStoppedPayload, McpToolCall,
McpToolDefinition, McpToolResult, McpToolsUpdatedPayload, ToolConverter,
};