MCP 支持两种通信传输方式:STDIO(标准输入/输出)或 SSE(服务器发送事件),两者均使用 JSON-RPC 2.0 格式的消息。STDIO 用于本地集成,而 SSE 用于基于网络的通信。
例如,如果我们希望在命令行中直接使用 MCP 服务,我们可以使用 STDIO 传输方法;如果我们希望在网页中使用 MCP 服务,我们可以使用 SSE 传输方法。
接下来,我们将为您开发一个基于 MCP 的智能商城服务助手,利用 SSE 类型的 MCP 服务,具有以下核心功能:
在这里我们使用 Anthropic Claude 3.5 Sonnet 模型作为 MCP 服务的 AI 助手,当然也可以选择其他支持工具调用的模型。
首先,需要一个产品微服务来暴露列出产品的 API。然后,提供一个订单微服务来暴露创建订单、检索库存信息等 API。
接下来的核心组件是中央 MCP SSE 服务器,它通过 SSE 协议将产品微服务和订单微服务的数据暴露给 LLM 作为工具。
最后,使用 MCP 客户端通过 SSE 协议连接到 MCP SSE 服务器,并与 LLM 互动。
接下来,我们将开始开发产品微服务和订单微服务,并暴露 API 接口。
首先,定义产品、库存和订单的类型。
// types/index.ts
export interface Product {
id: number;
name: string;
price: number;
description: string;
}
export interface Inventory {
productId: number;
quantity: number;
product?: Product;
}
export interface Order {
id: number;
customerName: string;
items: Array<{ productId: number; quantity: number }>;
totalAmount: number;
orderDate: string;
}
然后,我们可以使用 Express 来暴露产品微服务和订单微服务,并提供 API 接口。由于是模拟数据,我们将使用更简单的内存数据进行模拟,直接通过以下函数暴露数据。(在生产环境中,仍需实现带有数据库的微服务。)
// services/product-service.ts
import { Product, Inventory, Order } from "../types/index.js";
// 模拟数据存储
let products: Product[] = [
{
id: 1,
name: "智能手表Galaxy",
price: 1299,
description: "健康监测,运动追踪,支持多种应用",
},
{
id: 2,
name: "无线蓝牙耳机Pro",
price: 899,
description: "主动降噪,30小时续航,IPX7防水",
},
{
id: 3,
name: "便携式移动电源",
price: 299,
description: "20000mAh大容量,支持快充,轻薄设计",
},
{
id: 4,
name: "华为MateBook X Pro",
price: 1599,
description: "14.2英寸全面屏,3:2比例,100% sRGB色域",
},
];
// 模拟库存数据
let inventory: Inventory[] = [
{ productId: 1, quantity: 100 },
{ productId: 2, quantity: 50 },
{ productId: 3, quantity: 200 },
{ productId: 4, quantity: 150 },
];
let orders: Order[] = [];
export async function getProducts(): Promise<Product[]> {
return products;
}
export async function getInventory(): Promise<Inventory[]> {
return inventory.map((item) => {
const product = products.find((p) => p.id === item.productId);
return {
...item,
product,
};
});
}
export async function getOrders(): Promise<Order[]> {
return [...orders].sort(
(a, b) => new Date(b.orderDate).getTime() - new Date(a.orderDate).getTime()
);
}
export async function createPurchase(
customerName: string,
items: { productId: number; quantity: number }[]
): Promise<Order> {
if (!customerName || !items || items.length === 0) {
throw new Error("请求无效:缺少客户名称或商品");
}
let totalAmount = 0;
// 验证库存并计算总价
for (const item of items) {
const inventoryItem = inventory.find((i) => i.productId === item.productId);
const product = products.find((p) => p.id === item.productId);
if (!inventoryItem || !product) {
throw new Error(`商品ID ${item.productId} 不存在`);
}
if (inventoryItem.quantity < item.quantity) {
throw new Error(
`商品 ${product.name} 库存不足. 可用: ${inventoryItem.quantity}`
);
}
totalAmount += product.price * item.quantity;
}
// 创建订单
const order: Order = {
id: orders.length + 1,
customerName,
items,
totalAmount,
orderDate: new Date().toISOString(),
};
// 更新库存
items.forEach((item) => {
const inventoryItem = inventory.find(
(i) => i.productId === item.productId
)!;
inventoryItem.quantity -= item.quantity;
});
orders.push(order);
return order;
}
然后,我们可以使用 MCP 工具来暴露这些 API 接口,如下所示:
// mcp-server.ts
import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js";
import { z } from "zod";
import {
getProducts,
getInventory,
getOrders,
createPurchase,
} from "./services/product-service.js";
export const server = new McpServer({
name: "mcp-sse-demo",
version: "1.0.0",
description: "提供商品查询、库存管理和订单处理的MCP工具",
});
// 获取产品列表工具
server.tool("getProducts", "获取所有产品信息", {}, async () => {
console.log("获取产品列表");
const products = await getProducts();
return {
content: [
{
type: "text",
text: JSON.stringify(products),
},
],
};
});
// 获取库存信息工具
server.tool("getInventory", "获取所有产品的库存信息", {}, async () => {
console.log("获取库存信息");
const inventory = await getInventory();
return {
content: [
{
type: "text",
text: JSON.stringify(inventory),
},
],
};
});
// 获取订单列表工具
server.tool("getOrders", "获取所有订单信息", {}, async () => {
console.log("获取订单列表");
const orders = await getOrders();
return {
content: [
{
type: "text",
text: JSON.stringify(orders),
},
],
};
});
// 购买商品工具
server.tool(
"purchase",
"购买商品",
{
items: z
.array(
z.object({
productId: z.number().describe("商品ID"),
quantity: z.number().describe("购买数量"),
})
)
.describe("要购买的商品列表"),
customerName: z.string().describe("客户姓名"),
},
async ({ items, customerName }) => {
console.log("处理购买请求", { items, customerName });
try {
const order = await createPurchase(customerName, items);
return {
content: [
{
type: "text",
text: JSON.stringify(order),
},
],
};
} catch (error: any) {
return {
content: [
{
type: "text",
text: JSON.stringify({ error: error.message }),
},
],
};
}
}
);
在这里,我们总共定义了 4 个工具:
getProducts 获取所有产品信息getInventory 获取所有产品的库存信息getOrders 获取所有订单信息purchase 购买商品如果是 Stdio 类型的 MCP 服务,我们可以在命令行中直接使用这些工具,但现在我们需要使用 SSE 类型的 MCP 服务,所以我们还需要一个 MCP SSE 服务器来暴露这些工具。
接下来,我们将开始开发 MCP SSE 服务器,通过 SSE 协议将产品微服务和订单微服务的数据暴露为工具。
// mcp-sse-server.ts
import express, { Request, Response, NextFunction } from "express";
import cors from "cors";
import { SSEServerTransport } from "@modelcontextprotocol/sdk/server/sse.js";
import { server as mcpServer } from "./mcp-server.js"; // 重命名以避免命名冲突
const app = express();
app.use(
cors({
origin: process.env.ALLOWED_ORIGINS?.split(",") || "*",
methods: ["GET", "POST"],
allowedHeaders: ["Content-Type", "Authorization"],
})
);
// 存储活跃连接
const connections = new Map();
// 健康检查端点
app.get("/health", (req, res) => {
res.status(200).json({
status: "ok",
version: "1.0.0",
uptime: process.uptime(),
timestamp: new Date().toISOString(),
connections: connections.size,
});
});
// SSE 连接建立端点
app.get("/sse", async (req, res) => {
// 实例化SSE传输对象
const transport = new SSEServerTransport("/messages", res);
// 获取sessionId
const sessionId = transport.sessionId;
console.log(`[${new Date().toISOString()}] 新的SSE连接建立: ${sessionId}`);
// 注册连接
connections.set(sessionId, transport);
// 连接中断处理
req.on("close", () => {
console.log(`[${new Date().toISOString()}] SSE连接关闭: ${sessionId}`);
connections.delete(sessionId);
});
// 将传输对象与MCP服务器连接
await mcpServer.connect(transport);
console.log(`[${new Date().toISOString()}] MCP服务器连接成功: ${sessionId}`);
});
// 接收客户端消息的端点
app.post("/messages", async (req: Request, res: Response) => {
try {
console.log(`[${new Date().toISOString()}] 收到客户端消息:`, req.query);
const sessionId = req.query.sessionId as string;
// 查找对应的SSE连接并处理消息
if (connections.size > 0) {
const transport: SSEServerTransport = connections.get(
sessionId
) as SSEServerTransport;
// 使用transport处理消息
if (transport) {
await transport.handlePostMessage(req, res);
} else {
throw new Error("没有活跃的SSE连接");
}
} else {
throw new Error("没有活跃的SSE连接");
}
} catch (error: any) {
console.error(`[${new Date().toISOString()}] 处理客户端消息失败:`, error);
res.status(500).json({ error: "处理消息失败", message: error.message });
}
});
// 优雅关闭所有连接
async function closeAllConnections() {
console.log(
`[${new Date().toISOString()}] 关闭所有连接 (${connections.size}个)`
);
for (const [id, transport] of connections.entries()) {
try {
// 发送关闭事件
transport.res.write(
'event: server_shutdown\ndata: {"reason": "Server is shutting down"}\n\n'
);
transport.res.end();
console.log(`[${new Date().toISOString()}] 已关闭连接: ${id}`);
} catch (error) {
console.error(`[${new Date().toISOString()}] 关闭连接失败: ${id}`, error);
}
}
connections.clear();
}
// 错误处理
app.use((err: Error, req: Request, res: Response, next: NextFunction) => {
console.error(`[${new Date().toISOString()}] 未处理的异常:`, err);
res.status(500).json({ error: "服务器内部错误" });
});
// 优雅关闭
process.on("SIGTERM", async () => {
console.log(`[${new Date().toISOString()}] 接收到SIGTERM信号,准备关闭`);
await closeAllConnections();
server.close(() => {
console.log(`[${new Date().toISOString()}] 服务器已关闭`);
process.exit(0);
});
});
process.on("SIGINT", async () => {
console.log(`[${new Date().toISOString()}] 接收到SIGINT信号,准备关闭`);
await closeAllConnections();
process.exit( 0 );
});
// 启动服务器
const port = process.env.PORT || 8083;
const server = app.listen(port, () => {
console.log(
`[${new Date().toISOString()}] 智能商城 MCP SSE 服务器已启动,地址: http://localhost:${port}`
);
console.log(`- SSE 连接端点: http://localhost:${port}/sse`);
console.log(`- 消息处理端点: http://localhost:${port}/messages`);
console.log(`- 健康检查端点: http://localhost:${port}/health`);
});
在这里,我们使用 Express 暴露了一个 SSE 连接端点 /sse 用于接收客户端消息。使用 SSEServerTransport 创建一个 SSE 传输对象,并指定消息处理器端点为 /messages。
const transport = new SSEServerTransport("/messages", res);
传输对象创建后,我们可以将其连接到 MCP 服务器,如下所示:
// 将传输对象与MCP服务器连接
await mcpServer.connect(transport);
这样我们就可以通过 SSE 连接到端点 /sse 来接收客户端消息,并使用消息处理端点 /messages 来处理客户端消息。在接收客户端消息时,在 /messages 端点,我们需要使用 transport 对象来处理客户端消息:
// 使用transport处理消息
await transport.handlePostMessage(req, res);
这就是我们通常所说的列工具、调用工具等操作。
接下来,我们将开始开发 MCP 客户端,连接到 MCP SSE 服务器并通过 LLM 互动。对于客户端,我们可以开发命令行客户端或 Web 客户端。
我们已经介绍了命令行客户端,唯一的区别是我们现在需要使用 SSE 协议连接到 MCP SSE 服务器。
// 创建MCP客户端
const mcpClient = new McpClient({
name: "mcp-sse-demo",
version: "1.0.0",
});
// 创建SSE传输对象
const transport = new SSEClientTransport(new URL(config.mcp.serverUrl));
// 连接到MCP服务器
await mcpClient.connect(transport);
然后,其他操作与之前介绍的命令行客户端相同,涉及列出所有工具并将用户的提问连同工具一起发送给 LLM 进行处理。LLM 返回结果后,根据结果调用工具,将工具调用结果和历史消息发送给 LLM 进一步处理,从而获得最终结果。
对于 Web 客户端,其实与命令行客户端基本相同,只是我们现在将这些处理步骤封装在接口内,然后可以通过网页调用。
我们首先需要初始化 MCP 客户端,然后检索所有工具,将工具格式转换为 Anthropic 所需的数组形式,最后创建 Anthropic 客户端。
// 初始化MCP客户端
async function initMcpClient() {
if (mcpClient) return;
try {
console.log("正在连接到MCP服务器...");
mcpClient = new McpClient({
name: "mcp-client",
version: "1.0.0",
});
const transport = new SSEClientTransport(new URL(config.mcp.serverUrl));
await mcpClient.connect(transport);
const { tools } = await mcpClient.listTools();
// 转换工具格式为Anthropic所需的数组形式
anthropicTools = tools.map((tool: any) => {
return {
name: tool.name,
description: tool.description,
input_schema: tool.inputSchema,
};
});
// 创建Anthropic客户端
aiClient = createAnthropicClient(config);
console.log("MCP客户端和工具已初始化完成");
} catch (error) {
console.error("初始化MCP客户端失败:", error);
throw error;
}
}
接下来,我们根据自己的需求开发 API 接口。例如,我们在这里开发一个聊天接口来接收用户的问题,然后调用 MCP 客户端的工具,将工具调用结果连同历史消息发送给 LLM 进行处理,从而获得最终结果。代码如下:
// API: 聊天请求
apiRouter.post("/chat", async (req, res) => {
try {
const { message, history = [] } = req.body;
if (!message) {
console.warn("请求中消息为空");
return res.status(400).json({ error: "消息不能为空" });
}
// 构建消息历史
const messages = [...history, { role: "user", content: message }];
// 调用AI
const response = await aiClient.messages.create({
model: config.ai.defaultModel,
messages,
tools: anthropicTools,
max_tokens: 1000,
});
//