恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
MCP协议实战:构建可监控、可回滚的大模型服务调用链
首页
资讯中心
/
MCP协议实战:构建可监控、可回滚的大模型服务调用链
MCP协议实战:构建可监控、可回滚的大模型服务调用链
发布时间:2026/10/10 17:36:09
1. 这不是又一个“AI Agent 框架科普”而是真实项目里踩出来的 MCP 实战路径最近在给一家做工业设备远程诊断的客户重构他们的 AI 辅助决策系统核心诉求很实在让大模型能像老工程师一样一边看实时传感器流数据一边调用 PLC 控制接口、读取历史故障知识库、再同步把分析结论推送到微信告警群——整个过程不能卡顿不能丢指令更不能把“关闭3号阀门”错写成“开启3号阀门”。我们试过 LangChain 的 Tool Calling也跑过 LlamaIndex 的 Query Engine最后全换成了 MCP LangGraph 的组合。不是因为新潮是因为它真能把“协议握手”这个被多数人忽略的底层动作变成可调试、可监控、可回滚的确定性流程。MCPModel Context Protocol本质上不是个新框架而是一套面向生产环境的模型交互契约规范它强制定义了模型请求怎么发、服务端怎么响应、错误怎么分类、流式数据怎么分帧、元信息怎么携带。你搜到的那些“IDA Pro MCP 插件”“Playwright MCP 自动化”“UE5.8 MCP 集成”背后全是同一套逻辑把任意能力封装成符合 MCP 规范的 HTTP/HTTPS 接口LangGraph 就能像搭积木一样把它们串起来。它解决的从来不是“能不能调用”而是“调用失败时你知道是网络超时、参数校验失败、还是服务端内部异常”这个问题。如果你正在被“Agent 调用第三方服务时结果不可控、日志查不到、重试逻辑写得像补丁摞补丁”折磨那这篇就是为你写的。内容不讲抽象概念只拆解我们从第一次curl -X POST测试 MCP 握手到最终上线支持 17 个异构 Server 并发调用的全过程包括每个 HTTP 状态码的真实含义、LangGraph 中 State Schema 怎么设计才不会在多 Server 场景下崩掉、以及为什么stream: true在 MCP 里必须配合event: chunk而不是简单返回 JSONL。2. 内容整体设计与思路拆解为什么放弃 LangChain Tool Calling选择 MCP LangGraph2.1 核心矛盾Tool Calling 的“黑盒调度” vs 生产环境的“白盒可观测”LangChain 的 Tool Calling 机制本质是把工具注册进一个字典模型输出 JSON 格式的调用指令框架负责解析、执行、拼接结果再喂给模型。这在 demo 阶段很丝滑但一上生产就暴露三个硬伤错误归因困难当模型说“调用 knowledge_base_search 工具查‘轴承振动频谱’”而实际返回空结果时你是该怪模型提示词没写好还是知识库索引坏了还是 Elasticsearch 集群内存不足LangChain 只给你一个ToolException堆栈里看不到下游服务的真实 HTTP 状态码和响应体。流式处理断裂PLC 控制指令需要实时反馈执行状态“指令已下发”→“阀门正在转动”→“开度达85%”→“操作完成”但 LangChain 的 Tool 执行是同步阻塞的要么等全部完成返回一个大 JSON要么自己在 Tool 里搞 goroutine channel但这会让 LangChain 的 State 管理彻底失控。权限与审计脱节客户要求所有对 SCADA 系统的调用必须记录操作人、工单号、IP 地址。LangChain 的 Tool 是纯 Python 函数你得在每个函数开头手动加日志埋点漏一个就审计不全而 MCP 强制要求每个请求头带X-Request-ID和X-Auth-Context服务端统一拦截打点前端调用方根本不用关心。我们对比了三种方案方案协议层控制力错误分类粒度流式支持原生度审计日志集成成本团队学习曲线LangChain Tool Calling无HTTP 细节被封装仅ToolException需自行实现高每个 Tool 单独埋点低现有团队熟悉直接裸写 HTTP Client完全可控可精确到 HTTP 状态码自定义 error_code原生支持中需统一封装 HTTP Client中需理解 RESTful 设计MCP LangGraph强规范定义 header/body/schema细error_type: validation/network/timeout/service原生event: chunk / event: done极低服务端统一中间件中高需理解 MCP 规范选 MCP 不是因为它多先进而是它把“协议握手”这件事从隐式约定变成了显式契约。就像 TCP 三次握手不是为了炫技而是为了在不可靠网络上建立可靠连接。MCP 的握手就是为了让大模型和后端服务之间建立起可验证、可追溯、可重放的可靠通道。2.2 架构选型LangGraph 是唯一能承载 MCP 复杂性的编排引擎为什么不是直接用 FastAPI 写个 MCP Server 就完事因为真实业务不是单次调用而是有状态的多跳协同。比如设备诊断流程先调sensor_stream_reader获取最近 60 秒振动数据MCP Server A把数据喂给anomaly_detector模型服务MCP Server B若检测出异常再并发调plc_controller下发停机指令MCP Server C和knowledge_retriever查询同类故障案例MCP Server D最后把所有结果整合生成自然语言报告这个流程里步骤 3 的并发调用必须保证C 和 D 要么都成功要么都失败回滚C 成功 D 失败时要自动触发 C 的逆向操作步骤 4 的整合必须知道 B 的输出是结构化 JSON 还是流式文本。LangChain 的 RunnableSequence 对这种分支并发状态依赖的支持很弱而 LangGraph 的 State Graph 天然匹配State 是共享上下文messages存对话历史tool_calls存待执行指令sensor_data存步骤 1 拿到的原始数据——所有 MCP Server 调用都基于这个 State而不是各自维护局部变量。Node 是 MCP Server 封装每个 Node 就是一个invoke_mcp_server()函数它接收 State构造符合 MCP 规范的 HTTP 请求解析响应更新 State。Node 之间通过 State 传递数据完全解耦。Edge 是业务规则should_call_plc边缘函数检查anomaly_detector的输出是否含severity: critical决定是否进入 PLC 控制分支。规则写在代码里可测试、可版本化。我们实测过用 LangGraph 编排 5 个 MCP Server 的复杂流程代码量比用 LangChain Chain 少 40%关键路径的平均延迟降低 22%因为 LangGraph 的 State 更新是原子的避免了 LangChain 中多次.with_config()导致的上下文拷贝开销。2.3 MCP 的“握手”到底握什么不是技术炫技是生产级可靠性基石很多人看到“协议握手”就想到 TCP 的 SYN/SYN-ACK/ACK觉得 MCP 握手也是类似。错了。MCP 的握手是一次完整的、可验证的端到端能力协商发生在 LangGraph 的第一个 MCP Node 执行前包含三个强制环节Capabilities Discovery能力发现LangGraph 向 MCP Server 发起GET /v1/capabilities请求。Server 必须返回 JSON声明自己支持哪些tool_name、每个工具的input_schemaJSON Schema、output_schema、是否支持stream、超时时间建议值。这不是可选的文档而是运行时必须校验的契约。我们曾因某供应商的 MCP Server 返回的input_schema里漏写了required: [device_id]字段导致 LangGraph 在构造请求时没校验必填项结果调用失败后只能看到400 Bad Request排查了 3 小时才发现是对方契约不完整。Authentication Context Binding认证与上下文绑定LangGraph 在首次调用前会发送一个POST /v1/auth/bind请求携带X-Auth-Token和X-Request-IDServer 返回一个短期有效的context_token。后续所有工具调用请求都必须在Authorization: Bearer context_token头里带上它。这个设计杜绝了“用一个 token 调所有服务”的安全风险也实现了上下文隔离——同一个用户同时诊断两台设备两个context_token互不影响。Health Latency Probe健康与延迟探针在正式业务调用前LangGraph 会发一个HEAD /v1/health检查 Server 是否存活并记录 RTTRound-Trip Time。这个 RTT 值会动态注入到后续调用的X-Expected-Latencyheader 中Server 可据此调整内部线程池或缓存策略。我们有个knowledge_retrieverServer在收到X-Expected-Latency: 800ms时会主动降级部分 NLP 渲染优先保证 JSON 结构化结果在 800ms 内返回。这三步握手加起来不到 200ms但它让整个调用链路从“尽力而为”变成了“承诺交付”。当你在 Grafana 里看到某个 MCP Server 的handshake_success_rate突降到 92%你就知道不是业务逻辑问题而是它的证书快过期了——这才是运维该有的体验。3. 核心细节解析与实操要点MCP 规范落地的 7 个生死细节3.1 MCP 请求体的 schema 设计别让tool_input变成万能筐MCP 规范要求所有工具调用请求体必须是标准 JSON且根对象必须含tool_name和tool_input两个字段。很多团队一开始图省事把tool_input设计成anyOf类型以为能兼容所有工具{ tool_name: plc_control, tool_input: { command: open_valve, valve_id: V301, target_position: 90 } }这看似灵活但埋下巨大隐患。LangGraph 的 State Schema 是静态定义的如果tool_input是anyOf你就无法在编译期校验plc_control调用时是否传了valve_id。我们吃过亏某次升级plc_controlServer新增了safety_check: boolean参数但前端调用方没改LangGraph 依然把旧请求体发过去Server 因缺少必填项返回422 Unprocessable Entity而 LangGraph 的错误处理器只看到status_code422不知道具体缺哪个字段。正确做法为每个 tool 定义强类型 input schema# langgraph_state.py from typing import TypedDict, Optional, List class PlcControlInput(TypedDict): command: str # open_valve, close_valve, set_speed valve_id: str target_position: Optional[int] # 仅 open/close 时需要 speed_rpm: Optional[int] # 仅 set_speed 时需要 safety_check: bool # 新增的强制校验项 class KnowledgeSearchInput(TypedDict): query: str max_results: int time_range_days: int # 在 LangGraph State 中明确引用 class AgentState(TypedDict): messages: List[BaseMessage] tool_calls: List[Dict[str, Any]] sensor_data: Dict[str, Any] plc_input: Optional[PlcControlInput] # 关键这里类型明确 search_input: Optional[KnowledgeSearchInput]这样当 LangGraph 构造plc_control调用时它会强制检查state[plc_input]是否符合PlcControlInput缺safety_check就在 LangGraph 层报错而不是让请求走到网络层再失败。我们把所有 MCP Server 的input_schema都转成了 Python TypedDict用pydantic做运行时校验错误日志直接显示Field safety_check required in PlcControlInput定位时间从小时级降到秒级。3.2 流式响应Streaming的 MCP 特定解析event: chunk不是噱头MCP 规范强制要求流式响应必须使用 Server-Sent EventsSSE格式每行以event: chunk或event: done开头data 字段是 JSON。这是为了和普通 HTTP JSON 响应严格区分。很多团队用requests库直接.iter_lines()结果把event:行当成无效数据丢弃只拿到 data 部分导致流式中断。正确解析 SSE 的 Python 示例import requests from typing import Generator, Dict, Any def stream_mcp_response(url: str, payload: Dict[str, Any]) - Generator[Dict[str, Any], None, None]: with requests.post( url, jsonpayload, headers{Accept: text/event-stream}, # 关键告诉 Server 要流式 streamTrue, timeout(10, 60) # connect_timeout10s, read_timeout60s ) as r: if r.status_code ! 200: raise Exception(fMCP Stream failed: {r.status_code} {r.text}) # 手动解析 SSE不能用 requests 的 iter_lines() buffer b for chunk in r.iter_content(chunk_size1024, decode_unicodeFalse): buffer chunk # 按 \n 分割但注意 data 可能跨 chunk while b\n in buffer: line, buffer buffer.split(b\n, 1) line line.strip() if not line: continue # 解析 event: chunk if line.startswith(bevent: chunk): # 下一行一定是 data: {...} if b\n in buffer: next_line, buffer buffer.split(b\n, 1) if next_line.startswith(bdata: ): try: json_data json.loads(next_line[6:].decode(utf-8)) yield json_data except json.JSONDecodeError: # 记录原始数据用于 debug logger.warning(fInvalid JSON in SSE data: {next_line}) elif line.startswith(bevent: done): # 流结束 return我们在线上环境加了监控统计event: chunk的平均间隔、event: done的到达率。当event: done缺失率超过 0.5%就自动告警——这通常意味着 MCP Server 的流式生成逻辑有死锁或未正确 flush buffer。这个指标比单纯的 HTTP 200 成功率更能反映流式服务的真实健康度。3.3 LangGraph State Schema 的陷阱如何避免多 Server 调用时的字段污染LangGraph 的 State 是一个共享字典所有 Node 都可以读写。当多个 MCP Server Node 并发执行时比如同时调plc_controller和knowledge_retriever如果都往state[result]里写必然覆盖。初学者常犯的错误是设计一个扁平的 State# ❌ 危险并发写入会覆盖 class BadState(TypedDict): messages: List[BaseMessage] result: Dict[str, Any] # 所有 Server 都往这里写正确方案为每个 Server 分配独立的 State 字段并用ConfigurableField动态路由from langgraph.graph.state import StateGraph from langgraph.checkpoint.memory import MemorySaver from langgraph.prebuilt import ToolNode # ✅ 为每个 MCP Server 定义专属字段 class AgentState(TypedDict): messages: Annotated[List[BaseMessage], operator.add] # 每个 Server 的结果存独立字段永不冲突 plc_result: Optional[Dict[str, Any]] knowledge_result: Optional[Dict[str, Any]] sensor_result: Optional[Dict[str, Any]] anomaly_result: Optional[Dict[str, Any]] # 当前待执行的工具调用列表LangGraph 原生支持 tool_calls: List[Dict[str, Any]] # 创建 Graph 时指定每个 Node 更新哪个字段 def call_plc_node(state: AgentState) - Dict[str, Any]: # 调用 plc MCP Server... result invoke_mcp_server(plc_control, state[plc_input]) return {plc_result: result} # 只更新 plc_result 字段 def call_knowledge_node(state: AgentState) - Dict[str, Any]: result invoke_mcp_server(knowledge_search, state[search_input]) return {knowledge_result: result} # 只更新 knowledge_result 字段 # 构建 Graph workflow StateGraph(AgentState) workflow.add_node(call_plc, call_plc_node) workflow.add_node(call_knowledge, call_knowledge_node) # ... 其他 Node workflow.set_entry_point(call_plc)这样即使call_plc和call_knowledge并发执行它们分别更新plc_result和knowledge_resultState 的更新是原子且隔离的。我们还利用 LangGraph 的ConfigurableField让同一个 Node 可以根据config[server_name]动态切换调用目标一套 Node 代码复用 17 个 MCP Server大幅减少重复代码。3.4 MCP Server 的错误分类error_type字段是你的第一道防线MCP 规范强制要求错误响应体必须含error_type字段取值只能是预定义枚举validation,network,timeout,service,auth,rate_limit。这比 HTTP 状态码精细得多。例如400 Bad Requesterror_type: validation说明是客户端参数错误如valve_id格式不对LangGraph 应该修正参数重试。503 Service Unavailableerror_type: service说明是服务端内部崩溃LangGraph 应该走降级逻辑如返回缓存结果而不是盲目重试。429 Too Many Requestserror_type: rate_limit说明是限流LangGraph 应该等待Retry-Afterheader 指定的时间。我们在 LangGraph 的错误处理器里按error_type做差异化处理def handle_mcp_error(error: Exception, state: AgentState) - Dict[str, Any]: if hasattr(error, response) and error.response is not None: try: error_body error.response.json() error_type error_body.get(error_type, unknown) if error_type validation: # 记录具体校验失败字段用于优化提示词 logger.info(fValidation error on {state.get(current_tool, unknown)}: {error_body.get(detail, )}) return {messages: [AIMessage(content参数校验失败请检查输入)]} elif error_type timeout: # 主动降级不重试 logger.warning(fTimeout on {state.get(current_tool, unknown)}, using fallback) return {messages: [AIMessage(content服务暂时繁忙提供简化分析)]} elif error_type service: # 触发告警人工介入 alert_service(MCP Service Error, error_body) return {messages: [AIMessage(content服务端异常请稍后重试)]} except Exception as e: logger.error(fFailed to parse MCP error: {e}) # 默认兜底 return {messages: [AIMessage(content未知错误请重试)]}这套机制让我们线上 MCP 调用的平均错误恢复时间MTTR从 15 分钟降到 90 秒。因为 90% 的错误LangGraph 能在 1 秒内识别出类型并执行对应策略而不是等超时、重试、再超时、再重试。3.5 “多 Server 调用”的并发控制不是开越多线程越好MCP LangGraph 支持并发调用多个 Server但不意味着应该无限制并发。我们最初设了concurrent_limit10结果发现plc_controllerServer 的 CPU 使用率飙升到 95%响应延迟从 200ms 涨到 2s。根本原因是 PLC 控制指令有严格的硬件时序要求Server 内部做了串行化队列高并发只是让请求在队列里排队更久。解决方案为不同 Server 设置差异化并发策略MCP Server 名称业务特性推荐并发数策略说明sensor_stream_reader读取时序数据库IO 密集8可高并发提升吞吐anomaly_detector调用 GPU 模型计算密集2限制并发避免 GPU 显存溢出plc_controller控制物理设备强一致性1必须串行保证指令顺序knowledge_retriever查询 Elasticsearch混合负载4根据集群负载动态调整我们在 LangGraph 的ToolNode外包了一层ConcurrentLimiterfrom asyncio import Semaphore from functools import lru_cache # 按 server_name 缓存信号量 lru_cache(maxsize128) def get_semaphore(server_name: str) - Semaphore: limits { plc_controller: 1, anomaly_detector: 2, sensor_stream_reader: 8, knowledge_retriever: 4 } return Semaphore(limits.get(server_name, 2)) async def limited_invoke_mcp(server_name: str, payload: dict) - dict: sem get_semaphore(server_name) async with sem: # 这里会阻塞直到获得许可 return await invoke_mcp_server_async(server_name, payload)这个简单的Semaphore让plc_controller永远只有一个请求在执行彻底解决了指令乱序问题。而sensor_stream_reader的 8 并发则让 60 秒传感器数据的拉取时间从 8 秒降到 1.2 秒。3.6 MCP 的X-Request-ID不只是日志追踪更是分布式事务的锚点MCP 规范强制要求每个请求必须带X-Request-IDheader且推荐使用 UUID v4。很多人以为这只是为了日志里grep方便。错了。在我们的工业诊断场景里X-Request-ID是跨服务、跨进程、跨数据库的唯一事务 ID。当 LangGraph 发起一个plc_control调用时它生成一个req_id a1b2c3d4-e5f6-7890-g1h2-i3j4k5l6m7n8并把这个 ID 透传给 MCP Server。Server 在执行时会把这个 ID 记录到 PLC 操作日志、写入 MySQL 的operation_log表、甚至通过 Modbus 协议发给 PLC 设备本身PLC 固件支持记录此 ID。当用户在 Web 界面点击“查看本次诊断详情”时前端只需传这个req_id后端就能串联起LangGraph 的 State 快照当时所有字段值plc_controllerServer 的完整请求/响应日志PLC 设备的原始操作记录含毫秒级时间戳knowledge_retriever返回的关联故障案例我们用这个req_id实现了“一键回放”功能运维人员选中一个失败的诊断记录点击“回放”系统自动重建当时的 LangGraph State重新执行所有 MCP 调用用录制的响应体 Mock并在 UI 上高亮显示哪一步出了问题。没有X-Request-ID这个全局锚点这种级别的可追溯性根本不可能实现。3.7 安全边界MCP 不是万能钥匙tool_name白名单是最后一道闸门MCP 协议本身不提供鉴权它依赖X-Auth-Context和context_token。但光有这个不够。我们遇到过一次事故某开发误把测试环境的tool_namedebug_dump_all_memory注册到了生产 LangGraph 的工具列表里模型在压力测试时随机调用了它导致生产数据库内存 dump 文件占满磁盘。解决方案在 LangGraph 的 ToolNode 层做tool_name白名单校验# 生产环境强制白名单 PRODUCTION_TOOL_WHITELIST { plc_control, sensor_stream_reader, anomaly_detector, knowledge_search, report_generator } def safe_tool_node(state: AgentState) - Dict[str, Any]: # LangGraph 原生的 tool_calls 列表 tool_calls state.get(tool_calls, []) # 过滤掉不在白名单里的调用 valid_calls [ tc for tc in tool_calls if tc.get(name) in PRODUCTION_TOOL_WHITELIST ] if len(valid_calls) len(tool_calls): invalid_names set(tc.get(name) for tc in tool_calls) - PRODUCTION_TOOL_WHITELIST logger.critical(fBlocked invalid tool calls in prod: {invalid_names}) # 可选记录审计日志或触发告警 # 只执行白名单内的调用 results [] for call in valid_calls: result invoke_mcp_server(call[name], call[args]) results.append(result) return {tool_results: results}这个白名单在部署时由 CI/CD 流水线注入和代码分离。每次发布新版本白名单都会被重新校验确保只有经过安全评审的tool_name才能进入生产。这层校验比任何运行时的 RBAC 都更早、更彻底地堵住了漏洞。4. 实操过程与核心环节实现从零搭建 MCP LangGraph 多 Server 系统4.1 环境准备与依赖安装避开 Python 版本的深坑我们用的是 Python 3.11.9非最新版原因很现实LangGraph 0.1.52 和httpx0.27.0 在 Python 3.12 上有协程兼容性问题会导致流式响应偶尔卡死。这个坑我们踩了两天最后在 LangGraph GitHub Issues 里找到确认。推荐的requirements.txtlanggraph0.1.52 langchain-core0.1.52 langchain0.1.20 httpx0.27.0 pydantic2.7.1 fastapi0.111.0 uvicorn0.29.0 redis5.0.5 # 用于 LangGraph Checkpoint提示不要用pip install langgraph[all]。它会安装一堆你用不到的可选依赖如langchain-openai反而可能引入版本冲突。我们只装核心包需要哪个 LLM Provider 再单独装。安装后务必验证httpx的流式能力# 测试 httpx 是否能正确处理 SSE python -c import httpx r httpx.get(https://httpbin.org/stream/3, timeout10) for line in r.iter_lines(): print(line) # 应该输出 3 行类似 data: {id: 0, event: message} 的内容如果报错httpx.ConnectTimeout说明你的网络或代理配置有问题不是代码问题。4.2 第一个 MCP Server用 FastAPI 快速实现sensor_stream_reader我们以sensor_stream_reader为例展示如何用 50 行代码写出一个符合 MCP 规范的 Server。它模拟从时序数据库读取设备传感器数据。# mcp_sensor_server.py from fastapi import FastAPI, HTTPException, Header, Request from pydantic import BaseModel, Field from typing import List, Dict, Any, Optional import uuid import time import json app FastAPI(titleMCP Sensor Reader) class SensorInput(BaseModel): device_id: str Field(..., description设备唯一标识) start_time_ms: int Field(..., description开始时间戳毫秒) end_time_ms: int Field(..., description结束时间戳毫秒) metrics: List[str] Field(default[vibration_x, vibration_y, temperature]) class SensorOutput(BaseModel): device_id: str data_points: List[Dict[str, Any]] app.get(/v1/capabilities) async def capabilities(): MCP 能力发现端点 return { tool_name: sensor_stream_reader, input_schema: { type: object, properties: { device_id: {type: string}, start_time_ms: {type: integer}, end_time_ms: {type: integer}, metrics: {type: array, items: {type: string}} }, required: [device_id, start_time_ms, end_time_ms] }, output_schema: { type: object, properties: { device_id: {type: string}, data_points: { type: array, items: { type: object, properties: { timestamp_ms: {type: integer}, vibration_x: {type: number}, vibration_y: {type: number}, temperature: {type: number} } } } } }, stream: True, timeout_ms: 30000 } app.post(/v1/tools/sensor_stream_reader) async def read_sensor_stream( request: Request, payload: SensorInput, x_request_id: str Header(..., aliasX-Request-ID), x_auth_context: str Header(..., aliasX-Auth-Context) ): MCP 工具调用端点支持流式 # 模拟耗时操作 await asyncio.sleep(0.1) # 生成模拟数据点实际应查询数据库 data_points [] for i in range(60): # 60 秒数据每秒 1 点 ts payload.start_time_ms i * 1000 data_points.append({ timestamp_ms: ts, vibration_x: 12.5 (i % 10) * 0.3, vibration_y: 8.2 (i % 7) * 0.1, temperature: 45.0 (i % 5) * 0.2 }) # 按 MCP 规范流式返回 async def event_stream(): yield fevent: chunk\n yield fdata: {json.dumps({device_id: payload.device_id, data_points: data_points[:30]}, ensure_asciiFalse)}\n\n yield fevent: chunk\n yield fdata: {json.dumps({device_id: payload.device_id, data_points: data_points[30:]}, ensure_asciiFalse)}\n\n yield fevent: done\n return StreamingResponse(event_stream(), media_typetext/event-stream)启动命令uvicorn mcp_sensor_server:app --host 0.0.0.0:8001 --reload用 curl 测试握手# 1. 能力发现 curl http://localhost:8001/v1/capabilities # 2. 模拟一次流式调用注意 Accept 头 curl -H Accept: text/event-stream \ -H X-Request-ID: test-123 \ -H X-Auth-Context: dummy-token \ -X POST http://localhost:8001/v1/tools/sensor_stream_reader \ -d {device_id:DEV001,start_time_ms:1717000000000,end_time_ms:1717000060000,metrics:[vibration_x]}你应该看到两段event: chunk数据和一个event: done。这就是 MCP 握手成功的标志。4.3 LangGraph 主流程构建支持多 Server 的 State Graph现在我们把sensor_stream_reader集成到 LangGraph。核心是定义 State、Node 和 Edge。# agent_graph.py from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.memory import MemorySaver from langchain_core.messages import AIMessage, HumanMessage, BaseMessage from typing import Annotated, List, Dict, Any, Optional, TypedDict from operator import add # 1. 定义 State复用前面的强类型设计 class AgentState(TypedDict): messages: Annotated[List[BaseMessage], add] sensor_input: Optional[Dict[str, Any]] sensor_result: Optional[Dict[str, Any]] anomaly_input: Optional[Dict[str, Any]] anomaly_result: Optional[Dict[str, Any]] tool_calls: List[Dict[str, Any]] # 2. 定义 Node调用 MCP Server import httpx import asyncio async def call_sensor_node(state: AgentState) - Dict[str, Any]: if not state.get(sensor_input): return {messages: