原计划后置一个版本,经评估提前:迁移在新版启动时自动重写全部 客户端配置(settings.json/liangjing.json/全局注册),旧路由的实际 消费者只剩升级前的存量会话,而服务重启后它们本就断开重连。 单人环境下双路由过渡的保留价值趋近于零,一次收口减少一轮发版。 - tools_call 签名去 Option 组上下文,process_rpc 单参数化 - list_group_projects 的 group_id 改为 required 参数 - grep 'mcp/group' src-tauri/src/ 零命中,cargo test 93 全绿
201 lines
7.0 KiB
Rust
201 lines
7.0 KiB
Rust
use axum::{
|
||
extract::{Query, State},
|
||
http::StatusCode,
|
||
response::{
|
||
sse::{Event, KeepAlive, Sse},
|
||
IntoResponse, Response,
|
||
},
|
||
routing::{get, post},
|
||
Json, Router,
|
||
};
|
||
use serde::Deserialize;
|
||
use serde_json::{json, Value};
|
||
use std::{collections::HashMap, convert::Infallible, net::SocketAddr, sync::Arc};
|
||
use tokio::sync::{mpsc, Mutex};
|
||
use tokio_stream::{wrappers::ReceiverStream, StreamExt as _};
|
||
|
||
use crate::db;
|
||
|
||
/// SSE 会话表:session_id → 发送端
|
||
type Sessions = Arc<Mutex<HashMap<String, mpsc::Sender<(String, String)>>>>;
|
||
|
||
/// 从 settings 表读取 MCP 端口,默认 27190
|
||
pub fn read_port() -> u16 {
|
||
db::pool()
|
||
.get()
|
||
.ok()
|
||
.and_then(|c| {
|
||
c.query_row(
|
||
"SELECT value FROM settings WHERE key = 'mcp_port'",
|
||
[],
|
||
|r| r.get::<_, String>(0),
|
||
)
|
||
.ok()
|
||
})
|
||
.and_then(|v| v.parse().ok())
|
||
.unwrap_or(27190)
|
||
}
|
||
|
||
/// 启动 MCP HTTP 服务(同时支持 Streamable HTTP 和 SSE 两种传输)
|
||
pub async fn start(port: u16) {
|
||
let sessions: Sessions = Arc::new(Mutex::new(HashMap::new()));
|
||
|
||
let app = Router::new()
|
||
.route("/health", get(health))
|
||
// 单一全局服务器路由(lian-jing):寻址不依赖 group_id,工具显式传 project_id
|
||
.route("/mcp", post(handle_mcp).get(handle_sse_connect))
|
||
// SSE 消息端点:客户端通过此 URL 发送 JSON-RPC,响应经 SSE 推回
|
||
.route("/mcp/messages", post(handle_sse_message))
|
||
// Gitea webhook:接收 push/pull_request 事件写入飞轮
|
||
.route("/webhook/gitea/:project_id", post(crate::mcp_webhook::handle_gitea_webhook))
|
||
.with_state(sessions);
|
||
|
||
let addr = SocketAddr::from(([0, 0, 0, 0], port));
|
||
match tokio::net::TcpListener::bind(addr).await {
|
||
Ok(listener) => {
|
||
db::log_event("mcp_server", "info", &format!("MCP Server 已启动,端口 {port}(HTTP + SSE)"), None);
|
||
if let Err(e) = axum::serve(listener, app).await {
|
||
db::log_event("mcp_server", "error", &format!("Server 异常退出: {e}"), None);
|
||
}
|
||
}
|
||
Err(e) => db::log_event("mcp_server", "error", &format!("绑定端口 {port} 失败: {e}"), None),
|
||
}
|
||
}
|
||
|
||
async fn health() -> &'static str {
|
||
"炼境 MCP Server running"
|
||
}
|
||
|
||
// ── JSON-RPC 类型 ─────────────────────────────────────────────────────────────
|
||
|
||
#[derive(Deserialize)]
|
||
struct RpcRequest {
|
||
#[allow(dead_code)]
|
||
jsonrpc: String,
|
||
method: String,
|
||
id: Option<Value>,
|
||
#[serde(default)]
|
||
params: Option<Value>,
|
||
}
|
||
|
||
#[derive(Deserialize)]
|
||
struct SessionQuery {
|
||
session_id: String,
|
||
}
|
||
|
||
// ── 共享 RPC 处理逻辑 ─────────────────────────────────────────────────────────
|
||
|
||
async fn process_rpc(req: &RpcRequest) -> Result<Value, String> {
|
||
match req.method.as_str() {
|
||
"initialize" => Ok(json!({
|
||
"protocolVersion": "2024-11-05",
|
||
"capabilities": { "tools": {} },
|
||
"serverInfo": { "name": "炼境", "version": env!("CARGO_PKG_VERSION") }
|
||
})),
|
||
"tools/list" => Ok(crate::mcp_tools::tools_list_result()),
|
||
"tools/call" => Ok(crate::mcp_tools::tools_call(req.params.as_ref()).await),
|
||
_ => Err(format!("Method not found: {}", req.method)),
|
||
}
|
||
}
|
||
|
||
// ── Streamable HTTP(POST)─────────────────────────────────────────────────────
|
||
|
||
async fn handle_mcp(Json(req): Json<RpcRequest>) -> Response {
|
||
if req.id.is_none() {
|
||
return StatusCode::ACCEPTED.into_response();
|
||
}
|
||
let id = req.id.clone();
|
||
match process_rpc(&req).await {
|
||
Ok(result) => Json(json!({ "jsonrpc": "2.0", "id": id, "result": result })).into_response(),
|
||
Err(msg) => Json(json!({
|
||
"jsonrpc": "2.0", "id": id,
|
||
"error": { "code": -32601, "message": msg }
|
||
}))
|
||
.into_response(),
|
||
}
|
||
}
|
||
|
||
// ── SSE 传输(GET + POST /messages)─────────────────────────────────────────
|
||
|
||
/// GET /mcp — Claude Code SSE 握手,首条事件告知后续 POST URL
|
||
async fn handle_sse_connect(
|
||
State(sessions): State<Sessions>,
|
||
) -> Sse<impl futures_core::Stream<Item = Result<Event, Infallible>>> {
|
||
let session_id = uuid::Uuid::new_v4().to_string();
|
||
let (tx, rx) = mpsc::channel::<(String, String)>(32);
|
||
|
||
let port = read_port();
|
||
let endpoint_uri = format!(
|
||
"http://127.0.0.1:{}/mcp/messages?session_id={}",
|
||
port, session_id
|
||
);
|
||
let _ = tx.send(("endpoint".to_string(), endpoint_uri)).await;
|
||
sessions.lock().await.insert(session_id, tx);
|
||
|
||
let stream = ReceiverStream::new(rx)
|
||
.map(|(evt_name, data)| Ok::<Event, Infallible>(Event::default().event(evt_name).data(data)));
|
||
|
||
Sse::new(stream).keep_alive(KeepAlive::default())
|
||
}
|
||
|
||
#[cfg(test)]
|
||
mod tests {
|
||
use super::*;
|
||
|
||
fn rpc(method: &str) -> RpcRequest {
|
||
RpcRequest {
|
||
jsonrpc: "2.0".to_string(),
|
||
method: method.to_string(),
|
||
id: Some(json!(1)),
|
||
params: None,
|
||
}
|
||
}
|
||
|
||
/// 单一 lian-jing 服务器的握手路径:initialize / tools/list 必须可用
|
||
#[tokio::test]
|
||
async fn process_rpc_handshake() {
|
||
let init = process_rpc(&rpc("initialize")).await.unwrap();
|
||
assert_eq!(init["serverInfo"]["name"], "炼境");
|
||
|
||
let tools = process_rpc(&rpc("tools/list")).await.unwrap();
|
||
let list = tools["tools"].as_array().unwrap();
|
||
assert!(list.len() >= 20, "工具清单不应因去组化而缺失,实际 {}", list.len());
|
||
}
|
||
}
|
||
|
||
/// POST /mcp/messages?session_id=xxx
|
||
/// 立即返回 202,后台异步处理工具调用并经 SSE channel 推回响应
|
||
async fn handle_sse_message(
|
||
State(sessions): State<Sessions>,
|
||
Query(query): Query<SessionQuery>,
|
||
Json(req): Json<RpcRequest>,
|
||
) -> Response {
|
||
if req.id.is_none() {
|
||
return StatusCode::ACCEPTED.into_response();
|
||
}
|
||
|
||
let session_id = query.session_id.clone();
|
||
tokio::spawn(async move {
|
||
let id = req.id.clone();
|
||
let response_json = match process_rpc(&req).await {
|
||
Ok(result) => json!({ "jsonrpc": "2.0", "id": id, "result": result }),
|
||
Err(msg) => json!({
|
||
"jsonrpc": "2.0", "id": id,
|
||
"error": { "code": -32601, "message": msg }
|
||
}),
|
||
};
|
||
let mut guard = sessions.lock().await;
|
||
if let Some(tx) = guard.get(&session_id) {
|
||
if tx
|
||
.send(("message".to_string(), response_json.to_string()))
|
||
.await
|
||
.is_err()
|
||
{
|
||
guard.remove(&session_id);
|
||
}
|
||
}
|
||
});
|
||
|
||
StatusCode::ACCEPTED.into_response()
|
||
}
|