diff --git a/m10_mcp_basics/simple_client.py b/m10_mcp_basics/simple_client.py deleted file mode 100644 index c4096e8..0000000 --- a/m10_mcp_basics/simple_client.py +++ /dev/null @@ -1,19 +0,0 @@ -from mcp import ClientSession,StdioServerParameters -from mcp.client.stdio import stdio_client - -class SimpleClient: - def __init__(self,command:str,args:list[str],env:dict=None): - # 指定要启动的工具和参数 - self.params = StdioServerParameters(command=command,args=args,env=env) - - async def run_once(self,tool_name:str,tool_args:dict): - # 语法糖: async with 自动帮我们 打开连接 -> 运行 -> 关闭连接 - async with stdio_client(self.params) as (read,write): - # 建立父子进程管道(stdin/stdout) - async with ClientSession(read,write) as session: - # 用JSON-RPC与工具对话 - await session.initialize() - - # 直接调用工具 - result = await session.call_tool(tool_name,tool_args) - return result.content[0].text \ No newline at end of file diff --git a/m10_mcp_basics/simple_main.py b/m10_mcp_basics/simple_main.py deleted file mode 100644 index 5961566..0000000 --- a/m10_mcp_basics/simple_main.py +++ /dev/null @@ -1,24 +0,0 @@ -import asyncio -import os -from m10_mcp_basics.simple_client import SimpleClient -from config import AMAP_MAPS_API_KEY - -# 复制当前py进程的环境变量,并在复制的环境变量里新增一条,确保安全可控 -env_vars = os.environ.copy() -env_vars["AMAP_MAPS_API_KEY"] = AMAP_MAPS_API_KEY - -async def main(): - print('🔥 正在进行单次调用...') - client = SimpleClient( - command="npx", - args=["-y","@amap/amap-maps-mcp-server",AMAP_MAPS_API_KEY], - env=env_vars - ) - - # 这一步会经历:启动进程 - 握手 - 调用 - 杀进程 - result = await client.run_once("maps_text_search", {"keywords": "北京大学"}) - print(f'✅️ 结果:{result[:300]}') - - -if __name__ == "__main__": - asyncio.run(main()) \ No newline at end of file diff --git a/m11_mcp_advanced/README.md b/m11_mcp_advanced/README.md new file mode 100644 index 0000000..635bea6 --- /dev/null +++ b/m11_mcp_advanced/README.md @@ -0,0 +1,232 @@ +# 🧩 模块说明:MCP 高级篇 - 多模态协作协议客户端实现 + +> 📌 核心知识点:MCP协议高级应用|传输层封装|LangChain集成|流式输出|多服务管理 + +--- + +### 1. `agent_stream.py` (通用流式输出组件) + +实现通用的LangGraph事件流监听和可视化输出功能,提供友好的用户交互体验。 + +- ✅ 掌握点: + - LangGraph v2事件流的监听与处理 + - LLM流式吐字的实时渲染 + - 工具调用过程的可视化展示 + - 异步事件处理的最佳实践 + +- 功能: + - 监听LLM的流式输出并实时打印 + - 显示工具调用的开始和结束状态 + - 过滤内部包装工具,只显示自定义工具 + - 优化控制台输出格式,提升用户体验 + +> 💡 这是一个独立的工具组件,可以与任何LangGraph应用集成,用于增强用户交互体验。 + +--- + +### 2. `final_mcp_main.py` (官方库实现示例) + +使用官方`langchain_mcp_adapters`库实现的完整MCP应用示例,展示了如何快速集成MCP服务。 + +- ✅ 掌握点: + - 官方MultiServerMCPClient的使用方法 + - MCP服务的配置与初始化 + - LangGraph工作流的构建 + - 官方库与自定义组件的结合使用 + +- 功能演示: + - 初始化多服务器MCP客户端 + - 加载高德地图MCP服务 + - 构建基于LangGraph的地理位置助手 + - 使用自定义流式输出组件展示结果 + +> 💡 这是一个独立的示例应用,展示了如何使用官方库快速实现MCP功能,适合作为实际项目的参考。 + +--- + +### 组件系统:自定义MCP客户端实现 + +以下文件共同构成一个完整的自定义MCP客户端组件系统,实现了从传输层到应用层的完整封装。 + +--- + +### 3. `transports/base.py` (传输层协议接口) + +定义MCP传输层的抽象协议接口,为所有传输实现提供统一的规范。 + +- ✅ 掌握点: + - Python Protocol的使用方法 + - 抽象接口的设计原则 + - MCP协议的核心方法定义 + +- 功能: + - 定义MCP传输层必须实现的四个核心方法:connect、list_tools、call_tool、cleanup + - 提供类型注解,确保接口一致性 + - 为不同传输实现提供统一的调用方式 + +> 💡 这是整个组件系统的基础,定义了传输层的契约,使得上层代码可以与具体传输实现解耦。 + +--- + +### 4. `transports/http.py` (HTTP传输实现) + +实现基于HTTP协议的MCP传输层,支持与远程MCP服务器通信。 + +- ✅ 掌握点: + - HTTP JSON-RPC请求的实现 + - 异步HTTP客户端的使用 + - 会话管理与超时处理 + - 流式响应的处理 + +- 功能: + - 建立与远程MCP服务器的HTTP连接 + - 发送initialize请求并管理会话 + - 查询工具列表和调用工具 + - 处理普通JSON响应和SSE流式响应 + +> 💡 此实现支持远程MCP服务调用,适合构建分布式系统中的MCP客户端。 + +--- + +### 5. `transports/stdio.py` (标准输入输出传输实现) + +实现基于标准输入输出的MCP传输层,支持与本地MCP服务通信。 + +- ✅ 掌握点: + - AsyncExitStack资源管理 + - 子进程通信的实现 + - MCP协议的低级实现 + - 异步上下文管理器的应用 + +- 功能: + - 启动本地MCP服务进程 + - 建立标准输入输出管道通信 + - 管理MCP会话生命周期 + - 自动清理资源 + +> 💡 此实现支持本地MCP服务调用,适合开发和调试阶段使用。 + +--- + +### 6. `mcp_client.py` (客户端主类) + +实现MCP客户端的主类,封装传输层实现,提供统一的客户端接口。 + +- ✅ 掌握点: + - 工厂模式的应用 + - 依赖注入的实现 + - 客户端接口的设计 + - 错误处理的最佳实践 + +- 功能: + - 支持stdio和http两种传输方式 + - 封装连接、工具列表查询、工具调用和资源清理 + - 提供统一的客户端接口,隐藏传输层细节 + - 实现防御性编程,增强代码健壮性 + +> 💡 这是客户端组件的核心,为上层应用提供简洁易用的接口,同时屏蔽了底层传输的复杂性。 + +--- + +### 7. `mcp_bridge.py` (LangChain桥接器) + +实现MCP工具到LangChain工具的自动转换,使MCP服务能够无缝集成到LangChain生态中。 + +- ✅ 掌握点: + - JSON Schema到Pydantic模型的动态转换 + - LangChain工具的创建与配置 + - 批量工具加载的实现 + - 异步上下文管理器的应用 + +- 功能: + - 将MCP工具转换为LangChain可用的工具 + - 动态生成Pydantic参数模型 + - 支持批量加载多个MCP服务的工具 + - 管理MCP客户端的生命周期 + +> 💡 这是MCP与LangChain集成的关键组件,实现了两种生态系统之间的无缝对接。 + +--- + +### 8. `mcp_main.py` (完整应用示例) + +使用自定义MCP客户端组件实现的完整应用示例,展示了整个组件系统的协作使用。 + +- ✅ 掌握点: + - 组件系统的整体架构 + - 多MCP服务的配置与管理 + - LangGraph工作流的构建 + - 资源的统一管理 + +- 功能演示: + - 配置多个MCP服务(云端和本地) + - 批量加载MCP工具 + - 构建基于LangGraph的智能体 + - 使用流式输出展示结果 + +> 💡 这是整个组件系统的完整演示,展示了如何使用自定义实现构建功能完整的MCP应用。 + +--- + +### 组件系统架构图 + +``` +┌─────────────────────────────────────────────────────────┐ +│ 应用层 │ +│ ┌───────────────┐ ┌────────────────────────────────┐ │ +│ │ mcp_main.py │ │ final_mcp_main.py (官方库) │ │ +│ └───────────────┘ └────────────────────────────────┘ │ +│ │ │ │ +└──────────────┼─────────────────────┼────────────────────┘ + │ │ +┌──────────────┼─────────────────────┼────────────────────┐ +│ 集成层 │ +│ ┌───────────────┐ ┌─────────────────┐ │ +│ │ mcp_bridge.py│ │ agent_stream.py │ │ +│ └───────────────┘ └─────────────────┘ │ +│ │ │ +└──────────────┼──────────────────────────────────────────┘ + │ +┌──────────────┼──────────────────────────────────────────┐ +│ 客户端层 │ +│ ┌───────────────┐ │ +│ │ mcp_client.py│ │ +│ └───────────────┘ │ +│ │ │ +└──────────────┼──────────────────────────────────────────┘ + │ +┌──────────────┼──────────────────────────────────────────┐ +│ 传输层 │ +│ ┌───────────────┐ ┌───────────────┐ ┌─────────────┐ │ +│ │ transports/ │ │ transports/ │ │ transports/ │ │ +│ │ base.py │ │ http.py │ │ stdio.py │ │ +│ └───────────────┘ └───────────────┘ └─────────────┘ │ +└─────────────────────────────────────────────────────────┘ +``` + +--- + +### 🔔 全局注意事项 + +- **学习路径建议**: + 1. 先学习独立组件:`agent_stream.py` → `final_mcp_main.py` + 2. 再学习组件系统:`transports/base.py` → `transports/http.py` → `transports/stdio.py` → `mcp_client.py` → `mcp_bridge.py` → `mcp_main.py` + +- **环境准备**: + - 所有示例依赖根目录 `.env` 中的 API 密钥配置 + - MCP服务需要Node.js环境,确保已安装并配置正确路径 + - 运行前请确保已安装必要依赖:`pip install -r requirements.txt` + - 高德地图MCP服务需要 `AMAP_MAPS_API_KEY` 环境变量配置 + +- **运行说明**: + - 独立组件可以直接运行:`python final_mcp_main.py` + - 组件系统示例:`python mcp_main.py` + - 本地MCP服务需要先启动:`python -m m10_mcp_basics.streamable_http_server` + +--- + +### 💡 **扩展建议** +- 扩展MCP客户端,支持更多高级特性(如超时控制、重试机制等) +- 实现自定义的MCP服务,与客户端组件配合使用 +- 探索将MCP客户端与其他AI框架集成 +- 优化流式输出组件,支持更多展示效果 \ No newline at end of file diff --git a/m11_mcp_advanced/__init__.py b/m11_mcp_advanced/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/m10_mcp_basics/agent_stream.py b/m11_mcp_advanced/agent_stream.py similarity index 100% rename from m10_mcp_basics/agent_stream.py rename to m11_mcp_advanced/agent_stream.py diff --git a/m11_mcp_advanced/final_mcp_main.py b/m11_mcp_advanced/final_mcp_main.py new file mode 100644 index 0000000..19bc743 --- /dev/null +++ b/m11_mcp_advanced/final_mcp_main.py @@ -0,0 +1,107 @@ +import os +import asyncio + +# --- 核心:导入官方库 --- +from langchain_mcp_adapters.client import MultiServerMCPClient + +# LangChain/LangGraph 组件 +from langchain_openai import ChatOpenAI +from langchain_core.messages import SystemMessage +from langgraph.graph import StateGraph, MessagesState, START, END +from langgraph.prebuilt import ToolNode + +# 复用你的流式输出模块和配置 +from m11_mcp_advanced.agent_stream import run_agent_with_streaming +from config import OPENAI_API_KEY, AMAP_MAPS_API_KEY + +# === 配置 MCP 服务器 === +MCP_SERVERS = { + # 方式1.1: 云端代理 —— stdio模式 + "高德地图": { + "transport": "stdio", + "command": "npx", + "args": ["-y", "@amap/amap-maps-mcp-server"], + "env": {**os.environ, "AMAP_MAPS_API_KEY": AMAP_MAPS_API_KEY} + }, + + # 方式1.2: 云端MCP服务 —— Streamable HTTP模式 + # "高德地图" :{ + # "transport":"streamable_http", + # "url": f"https://mcp.amap.com/mcp?key={AMAP_MAPS_API_KEY}" + # }, + + # 方式2.1: 本地工具 —— stdio模式 + # "本地天气":{ + # "transport": "stdio", + # "command": "python", + # "args": ["-m", "m10_mcp_basics.stdio_server"], + # "env": None + # }, + + # 方式2.2:本地MCP服务 —— Streamable HTTP 模式 + # 注:此方法需要提前运行m10的 streamable_http_server.py + # "本地天气":{ + # "transport":"streamable_http", + # "url": "http://127.0.0.1:8001/mcp" + # } +} + + +def build_graph(available_tools): + """构建图逻辑 (保持不变)""" + if not available_tools: + print("⚠️ 未加载任何工具") + + llm = ChatOpenAI( + model="deepseek-chat", + api_key=OPENAI_API_KEY, + base_url="https://api.deepseek.com", + streaming=True + ) + + llm_with_tools = llm.bind_tools(available_tools) if available_tools else llm + + sys_prompt = "你是一个地理位置助手,请根据用户需求调用工具查询信息。" + + async def agent_node(state: MessagesState): + messages = [SystemMessage(content=sys_prompt)] + state["messages"] + return {"messages": [await llm_with_tools.ainvoke(messages)]} + + workflow = StateGraph(MessagesState) + workflow.add_node("agent", agent_node) + + if available_tools: + workflow.add_node("tools", ToolNode(available_tools)) + + def should_continue(state): + last_msg = state["messages"][-1] + return "tools" if last_msg.tool_calls else END + + workflow.add_edge(START, "agent") + workflow.add_conditional_edges("agent", should_continue) + workflow.add_edge("tools", "agent") + else: + workflow.add_edge(START, "agent") + workflow.add_edge("agent", END) + + return workflow.compile() + + +async def main(): + print("🔌 正在初始化 MCP 客户端...") + + client = MultiServerMCPClient(MCP_SERVERS) + + # 显式建立连接并获取工具 + # 注意:这个 client 对象会保持连接,直到脚本结束 + tools = await client.get_tools() + print(f"✅ 成功加载工具: {[t.name for t in tools]}") + + # 构建并运行 + app = build_graph(tools) + query = "帮我查一下杭州西湖附近的酒店" + await run_agent_with_streaming(app, query) + + +if __name__ == "__main__": + asyncio.run(main()) \ No newline at end of file diff --git a/m10_mcp_basics/mcp_bridge.py b/m11_mcp_advanced/mcp_bridge.py similarity index 68% rename from m10_mcp_basics/mcp_bridge.py rename to m11_mcp_advanced/mcp_bridge.py index 1a3968e..fd2ced9 100644 --- a/m10_mcp_basics/mcp_bridge.py +++ b/m11_mcp_advanced/mcp_bridge.py @@ -1,7 +1,9 @@ from typing import Dict,Any,Type from langchain_core.tools import StructuredTool -from m10_mcp_basics.mcp_client import MCPClient +from m11_mcp_advanced.mcp_client import MCPClient from pydantic import Field,create_model +from contextlib import AsyncExitStack + class LangChainMCPAdapter: """ @@ -61,7 +63,7 @@ class LangChainMCPAdapter: else: default_value = None - # 4.构建Pydantic字段定义 + # 4.构建Pydantic字段定义 —— create_model 要求的特定格式 fields[field_name] = (python_type,Field(default=default_value,description=description)) # 动态创建一个Pydantic模型类 @@ -95,4 +97,40 @@ class LangChainMCPAdapter: args_schema=args_model # 把说明书传给 LangChain ) langchain_tools.append(tool) - return langchain_tools \ No newline at end of file + return langchain_tools + + @classmethod + async def load_mcp_tools(cls,stack: AsyncExitStack, configs: list): + """ + 负责遍历配置,批量建立连接,收集所有工具。 + 使用stack将连接生命周期托管给上层 + """ + all_tools = [] + for conf in configs: + print(f'🔌 正在连接:{conf["name"]} == ({conf.get("transport","stdio")})...') + + # 根据 transport 类型创建不同的客户端 + transport = conf.get("transport","stdio") + if transport == "stdio": + # 初始化 Client + client = MCPClient( + transport="stdio", + command=conf["command"], + args=conf["args"], + env=conf.get("env") # 可选参数 + ) + else: # http + client = MCPClient( + transport="http", + url=conf["url"] + ) + + # 🔥:enter_async_context 替代了async with 缩进 + # 这样无论有多少个MCP,代码层级都不会变深 + adapter = await stack.enter_async_context(cls(client)) + # 批量获取一个MCP下的所有工具 + tools = await adapter.get_tools() + print(f' ✅️ 获取工具{[t.name for t in tools]}') + all_tools.extend(tools) + + return all_tools \ No newline at end of file diff --git a/m11_mcp_advanced/mcp_client.py b/m11_mcp_advanced/mcp_client.py new file mode 100644 index 0000000..3dfe313 --- /dev/null +++ b/m11_mcp_advanced/mcp_client.py @@ -0,0 +1,92 @@ +import uuid +from contextlib import AsyncExitStack +from typing import Optional,Literal + +from .transports.base import MCPTransport +from .transports.http import HttpMCPTransport +from .transports.stdio import StdioMCPTransport + + +class MCPClient: + + # 编辑器 _impl 必须满足MCPTransport协议 + _impl:MCPTransport + + def __init__( + self, + transport:Literal["stdio","http"]="stdio", + command:str=None, + args:list[str]=None, + env:dict=None, + url:str=None + ): + """ + MCP 客户端 - 支持 stdio 与 HTTP 两种传输方式 + + :param transport: 传输模式 "stdio" / "http" + :param command: stdio 模式的命令 (如npx) + :param args: stdio 模式的参数 + :param env: stdio 模式的环境变量 + :param url: http模式的端点 URL + """ + + if transport=="stdio": + if command is None: + raise ValueError("stdio 传输模式需要参数 'command'") + self._impl = StdioMCPTransport(command=command,args=args or [],env=env) + elif transport=="http": + if url is None: + raise ValueError("http 传输模式需要参数 'command'") + self._impl = HttpMCPTransport(url=url) + else: + raise ValueError(f"不支持的传输模式: {transport}") + + + + async def connect(self): + """ 建立MCP连接(stdio或HTTP) """ + await self._impl.connect() + + + async def list_tools(self): + """查询工具列表,为LLM建立上下文用""" + return await self._impl.list_tools() + + async def call_tool(self,name:str,args:dict): + """调用工具(工程化:加上防御性处理)""" + return await self._impl.call_tool(name,args) + + + async def cleanup(self): + """关闭MCP服务、会话和transport""" + return await self._impl.cleanup() + + + async def _http_request(self,method:str,params:dict=None): + """ 发送 HTTP JSON-RPC 请求 """ + payload = { + "jsonrpc":"2.0", + "id":str(uuid.uuid4()), + "method":method, + } + if params: + payload["params"] = params + + headers = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream" + } + if self.session_id: + headers["Mcp-Session-Id"] = self.session_id + + response = await self.http_client.post( + self.url, + json=payload, + headers=headers + ) + + if "Mcp-Session-Id" in response.headers: + self.session_id = response.headers["Mcp-Session-Id"] + + response.raise_for_status() + return response.json() \ No newline at end of file diff --git a/m10_mcp_basics/mcp_main.py b/m11_mcp_advanced/mcp_main.py similarity index 71% rename from m10_mcp_basics/mcp_main.py rename to m11_mcp_advanced/mcp_main.py index 5ab33e1..b6f2fb0 100644 --- a/m10_mcp_basics/mcp_main.py +++ b/m11_mcp_advanced/mcp_main.py @@ -9,26 +9,48 @@ from langgraph.graph import StateGraph,MessagesState,START,END from langgraph.prebuilt import ToolNode from config import OPENAI_API_KEY,AMAP_MAPS_API_KEY -from m10_mcp_basics.agent_stream import run_agent_with_streaming -from m10_mcp_basics.mcp_client import MCPClient -from m10_mcp_basics.mcp_bridge import LangChainMCPAdapter +from m11_mcp_advanced.agent_stream import run_agent_with_streaming +from m11_mcp_advanced.mcp_bridge import LangChainMCPAdapter # ===环境配置=== -# 环境兼容 -COMMAND = "npx.cmd" if sys.platform == "win32" else "npx" # 复制当前py进程的环境变量,并在复制的环境变量里新增一条,确保安全可控 env_vars = os.environ.copy() env_vars["AMAP_MAPS_API_KEY"] = AMAP_MAPS_API_KEY MCP_SERVER_CONFIGS = [ + # 方式1.1: 云端代理 —— stdio模式 { "name":"高德地图", # 打印使用了什么MCP,可移除 - "command":COMMAND, + "transport":"stdio", # 指定传输模式 + "command":"npx", "args":["-y", "@amap/amap-maps-mcp-server"], "env":env_vars } + + # 方式1.2: 云端MCP服务 —— Streamable HTTP模式 + # { + # "name":"高德地图", + # "transport":"http", + # "url": f"https://mcp.amap.com/mcp?key={AMAP_MAPS_API_KEY}" + # } + + # 方式2.1: 本地工具 —— stdio模式 + # { + # "name": "本地天气", + # "transport": "stdio", + # "command": "python", + # "args": ["-m", "m10_mcp_basics.stdio_server"], + # "env": None + # } + + # 方式2.2:本地MCP服务 —— Streamable HTTP 模式 + # { + # "name":"本地天气", + # "transport":"http", + # "url": "http://127.0.0.1:8001/mcp" + # } # {...} 之后MCP工具可随需求扩展增加 ] @@ -85,37 +107,13 @@ def build_graph(available_tools): return workflow.compile() -# ===MCP工具批量初始化=== -async def load_mcp_tools(stack:AsyncExitStack,configs:list): - """ - 负责遍历配置,批量建立连接,收集所有工具。 - 使用stack将连接生命周期托管给上层 - """ - all_tools = [] - for conf in configs: - print(f'🔌 正在连接:{conf["name"]}...') - # 初始化 Client - client = MCPClient( - command=conf["command"], - args=conf["args"], - env=conf.get("env") # 可选参数 - ) - # 🔥:enter_async_context 替代了async with 缩进 - # 这样无论有多少个MCP,代码层级都不会变深 - adapter = await stack.enter_async_context(LangChainMCPAdapter(client)) - # 批量获取一个MCP下的所有工具 - tools = await adapter.get_tools() - print(f' ✅️ 获取工具{[t.name for t in tools]}') - all_tools.extend(tools) - - return all_tools # ===主程序=== async def main(): # 使用ExitStack统一管理所有资源的关闭 async with AsyncExitStack() as stack: # A.插件(MCP)注入阶段 -- 允许为空 - dynamic_tools = await load_mcp_tools(stack,MCP_SERVER_CONFIGS) + dynamic_tools = await LangChainMCPAdapter.load_mcp_tools(stack,MCP_SERVER_CONFIGS) # B.图构建阶段 app = build_graph(available_tools=dynamic_tools) diff --git a/m11_mcp_advanced/transports/__init__.py b/m11_mcp_advanced/transports/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/m11_mcp_advanced/transports/base.py b/m11_mcp_advanced/transports/base.py new file mode 100644 index 0000000..e6e7b0d --- /dev/null +++ b/m11_mcp_advanced/transports/base.py @@ -0,0 +1,20 @@ +from typing import Protocol,List,Dict,Any + +class MCPTransport(Protocol): + """MCP 传输层协议 —— 所有 transport必须实现如下方法""" + + async def connect(self): + """建立连接""" + ... + + async def list_tools(self): + """获取工具列表""" + ... + + async def call_tool(self,name:str,args:dict): + """调用工具并返回文本结果""" + ... + + async def cleanup(self): + """清理资源""" + ... \ No newline at end of file diff --git a/m11_mcp_advanced/transports/http.py b/m11_mcp_advanced/transports/http.py new file mode 100644 index 0000000..4d820b1 --- /dev/null +++ b/m11_mcp_advanced/transports/http.py @@ -0,0 +1,130 @@ +import json +import uuid +from typing import Optional,Literal + +import httpx + + +class HttpMCPTransport: + def __init__(self,url:str=None): + """ + MCP 客户端 - 支持 stdio 与 HTTP 两种传输方式 + + :param url: http模式的端点 URL + """ + if not url: + raise ValueError("HTTP模式必须提供url参数") + self.url = url + self.session_id:Optional[str] = None + self.http_client:Optional[httpx.AsyncClient] = None + + + + async def connect(self): + """建立MCP长连接(一次连接,多次调用)""" + if self.http_client: + return + self.http_client = httpx.AsyncClient(timeout=60.0) + + # 发送 initialize 请求 + response = await self._http_request("initialize",{ + "protocolVersion":"2024-11-05", + "capabilities":{}, + "clientInfo":{"name":"mcp-client","version":"1.0"} + }) + + # 保存session ID + if response and "result" in response: + # Session ID在响应头内 + pass # 已在_http_request中处理 + + + async def list_tools(self): + """查询工具列表,为LLM建立上下文用""" + if not self.http_client: + raise RuntimeError("未连接,请先 connect()") + + result = await self._http_request("tools/list") + if result and "result" in result: + tools = result["result"].get("tools",[]) + return [ + { + "name":tool["name"], + "description":tool["description"], + "input_schema":tool.get("inputSchema",{}) + } + for tool in tools + ] + return [] + + + async def call_tool(self,name:str,args:dict): + """调用工具(工程化:加上防御性处理)""" + if not self.http_client: + raise RuntimeError("未连接,请先 connect()") + + result = await self._http_request("tools/call",{ + "name":name, + "arguments":args + }) + if result and "result" in result: + content = result["result"].get("content",{}) + if content and len(content): + return content[0].get("text",str(content[0])) + + return "工具执行成功,但无文本返回" + + + async def cleanup(self): + """关闭MCP服务、会话和transport""" + if self.http_client: + await self.http_client.aclose() + self.http_client = None + self.session_id = None + + + async def _http_request(self, method: str, params: dict = None): + """ 发送 HTTP JSON-RPC 请求 """ + payload = { + "jsonrpc": "2.0", + "id": str(uuid.uuid4()), + "method": method, + } + if params: + payload["params"] = params + + headers = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream" + } + if self.session_id: + headers["Mcp-Session-Id"] = self.session_id + + response = await self.http_client.post( + self.url, + json=payload, + headers=headers + ) + + if "Mcp-Session-Id" in response.headers: + self.session_id = response.headers["Mcp-Session-Id"] + + response.raise_for_status() + + # 判断响应类型 + content_type = response.headers.get("Content-Type", "") + + if "text/event-stream" in content_type: + # SSE 流式响应(本地服务器) + for line in response.text.split('\n'): + line = line.strip() + if line.startswith("data:"): + data_str = line[5:].strip() + try: + return json.loads(data_str) + except json.JSONDecodeError: + pass + return None + else: + # 普通 JSON 响应(云端服务) + return response.json() diff --git a/m10_mcp_basics/mcp_client.py b/m11_mcp_advanced/transports/stdio.py similarity index 85% rename from m10_mcp_basics/mcp_client.py rename to m11_mcp_advanced/transports/stdio.py index bb5da83..0e488c0 100644 --- a/m10_mcp_basics/mcp_client.py +++ b/m11_mcp_advanced/transports/stdio.py @@ -1,11 +1,20 @@ from contextlib import AsyncExitStack -from typing import Optional +from typing import Optional,Literal + from mcp import ClientSession,StdioServerParameters from mcp.client.stdio import stdio_client -class MCPClient: - def __init__(self,command:str,args:list[str],env:dict=None): +class StdioMCPTransport: + def __init__(self,command:str=None,args:list[str]=None,env:dict=None): + """ + MCP 客户端 - 支持 stdio 与 HTTP 两种传输方式 + + :param transport: 传输模式 "stdio" / "http" + :param command: stdio 模式的命令 (如npx) + :param args: stdio 模式的参数 + :param env: stdio 模式的环境变量 + """ # MCP启动方式(npx/uvx/python -m xxx) self.params = StdioServerParameters(command=command,args=args,env=env) # 工程核心:资源栈 @@ -13,11 +22,12 @@ class MCPClient: # 连接会话(长连接) self.session:Optional[ClientSession]=None + + async def connect(self): """建立MCP长连接(一次连接,多次调用)""" if self.session: return # 已连接无需重复 - # 进入transport(读/写管道) transport = await self.exit_stack.enter_async_context( stdio_client(self.params) @@ -29,6 +39,8 @@ class MCPClient: # 等待MCP服务器返回工具清单 await self.session.initialize() + + async def list_tools(self): """查询工具列表,为LLM建立上下文用""" if not self.session: @@ -53,6 +65,8 @@ class MCPClient: for tool in result.tools ] + + async def call_tool(self,name:str,args:dict): """调用工具(工程化:加上防御性处理)""" if not self.session: @@ -66,6 +80,8 @@ class MCPClient: return "工具执行成功,但无文本返回" + + async def cleanup(self): """关闭MCP服务、会话和transport""" if self.session: