ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

构建AI Agent统一发现层:ARD架构原理与Python实战

构建AI Agent统一发现层:ARD架构原理与Python实战 1. 从“信息孤岛”到“统一发现”为什么我们需要 ARD如果你最近在折腾 AI Agent 或者 MCPModel Context Protocol开发大概率会遇到一个让人头疼的场景你手头有好几个 Agent每个 Agent 又可能连接着不同的 MCP 服务器比如一个负责搜索网络信息一个负责读取数据库还有一个能操作你的本地文件系统。当你想让一个主 Agent 去完成一个复杂任务时你发现你得像一个“人工调度员”一样在配置文件里手动写下每个 Agent 的地址、端口、支持的能力。更麻烦的是一旦某个 Agent 的地址变了或者新加入了一个更强大的工具 Agent你就得去更新所有相关的配置。整个系统就像一堆各自为政的“信息孤岛”彼此知道对方存在但缺乏一个高效的“通讯录”和“广播系统”。这就是ARDAgent Resource Discovery要解决的核心问题。ARD即 Agent 资源统一发现其目标就是为分布式的 AI Agent 和 MCP 服务器建立一个动态的、可扩展的“服务注册与发现”层。简单来说它想让 Agent 像在局域网里发现打印机一样能被自动“搜到”并且能清晰地知道彼此能“干什么”。这不仅仅是技术上的便利更是构建复杂、可协作的 Agent 生态系统的基础设施。没有它多 Agent 协作要么停留在简单、静态的配置阶段要么就需要极其复杂的手工编排难以实现真正的智能化和自动化。最近随着 Claude Code、Cursor 等智能编码工具对 MCP 协议的支持日益成熟以及各类搜索、绘图、数据分析等专用 MCP 服务器的涌现如何高效地管理和集成这些资源成了一个显性的痛点。你或许在社区里看到过这样的问题“如何把 Tavily 搜索 MCP 服务器添加到我的 Code 环境中” 或者 “为什么我的 Figma MCP 插件还原度很低”。这些问题背后或多或少都指向了资源发现、配置和理解的缺失。ARD 正是为了系统性地解决这类问题而提出的架构思路。2. 理解 MCP 与 A2AARD 的两大基石要搞懂 ARD 在做什么必须先厘清它所要管理的两个核心实体MCPModel Context Protocol和A2AAgent-to-Agent通信。它们是现代 AI 应用架构中两种主流的资源交互模式。2.1 MCP为 LLM 扩展能力的“标准插座”MCP 协议你可以把它想象成给大语言模型LLM用的“USB-C 标准”。在没有 MCP 之前每个工具、每个数据库想被 LLM 调用都需要开发一套独立的、五花八门的集成方式就像老式手机各有各的充电口。MCP 定义了一套标准协议任何符合 MCP 的服务MCP Server都可以通过一个标准的“插座”MCP Client暴露自己的能力给 LLM。一个典型的 MCP 服务器会向客户端声明工具Tools我能执行哪些操作比如search_web,read_file。资源Resources我能提供哪些数据源比如file:///path/to/doc.md。提示词模板Prompts我有哪些预制好的对话模板例如一个tavily-mcp服务器会声明一个web_search工具一个filesystem-mcp服务器会声明读写文件工具和文件资源。你的 AI 编码助手如 Claude Code作为 MCP 客户端加载这些服务器后LLM 就能直接调用这些工具仿佛它们是天生的能力一样。那么MCP 如何与 ARD 关联在 ARD 的愿景中一个 MCP 服务器启动后可以向 ARD 服务注册自己“嗨我是一个 MCP 服务器我的地址是tcp://192.168.1.100:8000我提供了web_search和get_news这两个工具。” 这样其他需要搜索能力的 Agent 或客户端不需要预先知道这个服务器的存在和地址只需向 ARD 查询“谁能提供搜索服务”就能自动发现并连接上它。2.2 A2A智能体间的直接对话与协作如果说 MCP 是 LLM 与“工具世界”的桥梁那么 A2A 就是“智能体世界”内部的语言。A2A 通信指的是不同的 AI Agent 之间为了完成共同目标而进行的直接信息交换与任务协作。例如你可能设计了一个“规划 Agent”它擅长分解复杂任务一个“执行 Agent”它擅长调用具体工具还有一个“验证 Agent”负责检查结果质量。它们之间需要通过 A2A 协议来传递任务描述、执行结果、状态同步等信息。与 MCP 主要面向“工具调用”不同A2A 更侧重于“会话”与“协作”消息格式可能更灵活包含更多的上下文、意图和协商内容。A2A 在 ARD 中的角色ARD 需要能够发现这些 Agent 本身。一个 Agent 启动后可以向 ARD 注册“我是一个 AgentID 是planner_agent_01我擅长任务分解我的通信端点endpoint是grpc://10.0.0.5:50051我支持TaskDecomposition和Coordination这两种能力协议。” 当“执行 Agent”需要找一个“规划者”时它就可以通过 ARD 发现并联系上这个规划 Agent。MCP 与 A2A 的融合场景一个强大的 Agent 本身既可以作为 A2A 网络中的一个节点也可以同时作为一个 MCP 服务器对外提供工具化能力。ARD 的统一发现机制正是要管理这种混合的、多维度的资源网络。3. ARD 的核心架构与工作原理ARD 并非一个单一的软件而是一套架构模式和协议规范。一个典型的 ARD 系统通常包含以下几个核心组件其工作流程可以类比为互联网中的 DNS域名系统加上服务网格Service Mesh的注册中心。3.1 核心组件解析资源提供者Resource Provider角色希望被发现的实体即 MCP 服务器或 AI Agent。职责启动时向ARD 注册中心发送注册请求。请求中需要包含关键的元数据Metadata。关键元数据示例身份标识唯一的 ID如 UUID、名称。访问端点如何连接它如http://host:port,tcp://host:port,grpc://host:port。能力描述对于 MCP 服务器提供的工具列表、资源模式、提示词模板。对于 Agent支持的任务类型、通信协议如基于 OpenAI 函数调用、自定义 gRPC、技能Skills描述。健康状态定期发送心跳heartbeat以表明自己在线。元数据标签用于分类和筛选的键值对如type: “mcp-server”,category: “search”,version: “1.2.0”。ARD 注册中心Registry角色系统的核心一个持久化的服务可以是中心化的也可以是去中心化的如 Gossip 协议。职责接收并存储资源提供者的注册信息。维护资源的实时健康状态通过心跳机制自动剔除失效节点。提供查询接口让资源消费者能根据条件查找资源。技术选型参考在实践中可以基于成熟的协调服务实现如etcd、Consul、ZooKeeper或者专门为微服务设计的产品如Nacos。对于轻量级或边缘场景也可能采用简单的 HTTP 注册表。资源消费者Resource Consumer角色需要发现和使用资源的实体通常是另一个 Agent、MCP 客户端或编排引擎。职责向 ARD 注册中心发起查询例如“查找所有类型为mcp-server且标签categorysearch的资源”。获取到资源列表后根据端点信息建立连接并调用其能力。可选ARD 客户端 SDK为了简化集成通常会提供不同语言的 SDK。这个 SDK 封装了注册、发现、心跳等底层通信细节让开发者在自己的 MCP 服务器或 Agent 中只需几行代码就能接入 ARD 网络。3.2 工作流程与交互协议一个完整的工作流程通常遵循以下步骤我们可以通过一个“数据分析 Agent 需要搜索功能”的场景来具体说明注册阶段tavily-mcp-server启动加载自身的工具定义web_search。它通过 ARD SDK向部署在http://ard-registry:8080的注册中心发送注册请求。请求体包含{ “id”: “tavily-001”, “name”: “Tavily Web Search”, “type”: “mcp-server”, “endpoint”: “tcp://tavily-host:8000”, “tools”: [“web_search”], “metadata”: {“category”: “search”, “provider”: “tavily”} }。注册中心将其存入资源目录并开始期待它的心跳。发现阶段># ard_registry.py import asyncio import uuid from datetime import datetime, timedelta from typing import Dict, List, Optional from fastapi import FastAPI, HTTPException, BackgroundTasks from fastapi.middleware.cors import CORSMiddleware from pydantic import BaseModel, Field from contextlib import asynccontextmanager # 定义资源模型 class ResourceMetadata(BaseModel): category: Optional[str] None provider: Optional[str] None tags: Dict[str, str] {} class Resource(BaseModel): id: str Field(default_factorylambda: str(uuid.uuid4())) name: str type: str # e.g., mcp-server, agent endpoint: str # e.g., tcp://host:port, http://host:port # 对于MCP服务器可存储其声明的工具 tools: List[str] [] # 对于Agent可存储其技能 skills: List[str] [] metadata: ResourceMetadata Field(default_factoryResourceMetadata) last_heartbeat: datetime Field(default_factorydatetime.utcnow) is_healthy: bool True class RegisterRequest(BaseModel): name: str type: str endpoint: str tools: List[str] [] skills: List[str] [] metadata: ResourceMetadata Field(default_factoryResourceMetadata) class QueryRequest(BaseModel): type: Optional[str] None metadata_filter: Optional[Dict[str, str]] None # 全局存储和任务 resources: Dict[str, Resource] {} subscriptions: Dict[str, asyncio.Queue] {} # 用于服务发现的订阅队列 async def health_checker(): 后台任务定期检查资源健康状态 while True: await asyncio.sleep(30) # 每30秒检查一次 now datetime.utcnow() to_remove [] for rid, resource in resources.items(): # 如果超过90秒未收到心跳标记为不健康 if now - resource.last_heartbeat timedelta(seconds90): resource.is_healthy False # 可选一段时间后彻底删除 # if now - resource.last_heartbeat timedelta(minutes5): # to_remove.append(rid) _notify_subscribers(update, resource) elif not resource.is_healthy: # 恢复健康如果心跳恢复这个逻辑需要在心跳接口处理 pass for rid in to_remove: resources.pop(rid, None) _notify_subscribers(delete, rid) def _notify_subscribers(event: str, data): 通知所有订阅者资源变更 for queue in subscriptions.values(): try: queue.put_nowait({event: event, data: data}) except asyncio.QueueFull: pass asynccontextmanager async def lifespan(app: FastAPI): # 启动时运行健康检查任务 task asyncio.create_task(health_checker()) yield # 关闭时取消任务 task.cancel() app FastAPI(lifespanlifespan) app.add_middleware(CORSMiddleware, allow_origins[*], allow_methods[*], allow_headers[*]) app.post(/register) async def register(resource_req: RegisterRequest): 资源注册接口 resource Resource(**resource_req.dict()) resources[resource.id] resource _notify_subscribers(create, resource) return {id: resource.id, message: Registered successfully} app.post(/heartbeat/{resource_id}) async def heartbeat(resource_id: str): 资源心跳接口 if resource_id not in resources: raise HTTPException(status_code404, detailResource not found) resource resources[resource_id] resource.last_heartbeat datetime.utcnow() if not resource.is_healthy: resource.is_healthy True _notify_subscribers(update, resource) return {status: ok} app.post(/discover) async def discover(query: QueryRequest): 资源发现接口一次性查询 matched [] for resource in resources.values(): if not resource.is_healthy: continue if query.type and resource.type ! query.type: continue if query.metadata_filter: meta resource.metadata.dict() # 简单匹配查询的键值对必须完全匹配资源的metadata if not all(meta.get(k) v for k, v in query.metadata_filter.items()): continue matched.append(resource) return {resources: matched} app.get(/watch) async def watch_resources(background_tasks: BackgroundTasks): 服务发现的长连接接口SSE用于实时监听资源变化 queue asyncio.Queue() subscription_id str(uuid.uuid4()) subscriptions[subscription_id] queue async def event_generator(): try: # 首先发送全量资源 for resource in resources.values(): if resource.is_healthy: yield fdata: {resource.json()}\n\n while True: event_data await queue.get() # 根据事件类型格式化数据 yield fevent: {event_data[event]}\ndata: {event_data[data]}\n\n except asyncio.CancelledError: pass finally: subscriptions.pop(subscription_id, None) return EventSourceResponse(event_generator()) if __name__ __main__: import uvicorn uvicorn.run(app, host0.0.0.0, port8080)这个注册中心提供了四个核心接口POST /register: 用于资源注册。POST /heartbeat/{id}: 用于上报心跳。POST /discover: 用于一次性查询资源。GET /watch: 一个 SSE 端点客户端可以订阅资源变更事件创建、更新、删除实现实时发现。4.2 实现一个 MCP 服务器并接入 ARD接下来我们实现一个最简单的“模拟搜索” MCP 服务器并在启动时自动注册到我们的 ARD 中心。# mcp_search_server.py import asyncio import json from mcp import ClientSession, StdioServerParameters from mcp.server import Server from mcp.server.models import InitializationOptions import httpx import logging # 配置 ARD 注册中心地址 ARD_REGISTRY_URL http://localhost:8080 RESOURCE_NAME Simulated-Search-Server RESOURCE_TYPE mcp-server async def register_with_ard(): 向 ARD 注册中心注册本 MCP 服务器 registration_data { name: RESOURCE_NAME, type: RESOURCE_TYPE, endpoint: stdio, # 对于Stdio服务器端点可能是一个启动命令这里简化为标识 # 在实际中可能需要一个网络端点。这里我们假设ARD能通过进程名或SSE发现实际连接方式。 # 更完善的实现服务器应打开一个网络端口并将 endpoint 设为 tcp://host:port tools: [simulated_search], metadata: { category: search, provider: demo, description: A demo MCP server that simulates web search. } } async with httpx.AsyncClient() as client: try: resp await client.post(f{ARD_REGISTRY_URL}/register, jsonregistration_data, timeout5.0) resp.raise_for_status() result resp.json() resource_id result[id] logging.info(fSuccessfully registered with ARD, resource ID: {resource_id}) return resource_id except Exception as e: logging.error(fFailed to register with ARD: {e}) return None async def send_heartbeat(resource_id: str): 定期向 ARD 发送心跳 async with httpx.AsyncClient() as client: while True: await asyncio.sleep(30) # 每30秒发送一次 try: await client.post(f{ARD_REGISTRY_URL}/heartbeat/{resource_id}, timeout3.0) except Exception as e: logging.warning(fHeartbeat failed: {e}) async def main(): logging.basicConfig(levellogging.INFO) # 1. 启动时向 ARD 注册 resource_id await register_with_ard() if resource_id: # 2. 启动心跳任务 heartbeat_task asyncio.create_task(send_heartbeat(resource_id)) # 3. 创建并运行 MCP 服务器 server Server(simulated-search-server) server.list_tools() async def handle_list_tools(): return [ { name: simulated_search, description: Simulates a web search and returns mock results., inputSchema: { type: object, properties: { query: {type: string, description: The search query.} }, required: [query] } } ] server.call_tool() async def handle_call_tool(name: str, arguments: dict): if name simulated_search: query arguments.get(query, ) # 模拟搜索逻辑 results [ {title: fResult 1 about {query}, snippet: This is a simulated snippet 1., url: http://example.com/1}, {title: fResult 2 about {query}, snippet: This is a simulated snippet 2., url: http://example.com/2}, ] return { content: [{type: text, text: json.dumps(results, indent2)}] } raise ValueError(fUnknown tool: {name}) # 使用 Stdio 传输这是 MCP 常见方式被 Claude Code、Cursor 等支持 async with server.run_stdio() as (read_stream, write_stream): session ClientSession(read_stream, write_stream) await session.initialize(InitializationOptions(root_namespaceNone)) # 主循环处理来自客户端的请求 async for message in session.listen(): # 服务器逻辑已由 server 装饰器处理 pass # 4. 服务器关闭时应发送注销请求此处省略可通过 signal 处理 # if resource_id: # await deregister_from_ard(resource_id) if __name__ __main__: asyncio.run(main())这个 MCP 服务器做了三件事启动时向 ARD 注册中心注册自己声明其工具和能力。启动一个后台任务定期发送心跳以保持健康状态。运行一个符合 MCP 标准的 Stdio 服务器提供simulated_search工具。4.3 实现一个能发现并使用资源的智能 Agent最后我们创建一个简单的 Agent它会主动从 ARD 发现搜索服务并使用它。# smart_agent.py import asyncio import httpx import json from mcp import ClientSession, StdioServerParameters import subprocess import logging ARD_REGISTRY_URL http://localhost:8080 async def discover_mcp_servers(category: str search): 从 ARD 发现指定类别的 MCP 服务器 query { type: mcp-server, metadata_filter: {category: category} } async with httpx.AsyncClient() as client: try: resp await client.post(f{ARD_REGISTRY_URL}/discover, jsonquery, timeout5.0) resp.raise_for_status() result resp.json() return result.get(resources, []) except Exception as e: logging.error(fDiscovery failed: {e}) return [] async def connect_and_use_mcp_server(server_info: dict): 连接到一个 MCP 服务器并使用其工具 # 注意这里 server_info[endpoint] 是 ‘stdio’。在实际中我们需要知道如何启动这个服务器。 # 这是一个简化示例。更真实的场景是endpoint 可能是一个网络地址或者附带启动命令。 # 假设我们通过一个已知的启动命令来连接这个模拟搜索服务器。 server_params StdioServerParameters( commandpython, args[mcp_search_server.py] # 这里需要能启动目标服务器的命令 ) async with ClientSession(server_params) as session: await session.initialize() # 列出服务器工具 tools await session.list_tools() logging.info(fConnected to server. Available tools: {[t.name for t in tools.tools]}) # 使用搜索工具 for tool in tools.tools: if search in tool.name: result await session.call_tool(tool.name, arguments{query: AI Agent frameworks 2024}) content result.content if content: print(fSearch Results from {server_info[name]}:) print(content[0].text) break async def main(): logging.basicConfig(levellogging.INFO) # 1. 发现搜索服务 logging.info(Discovering search MCP servers from ARD...) servers await discover_mcp_servers(search) if not servers: logging.warning(No search servers found.) return logging.info(fFound {len(servers)} server(s).) # 2. 连接并使用第一个找到的服务器在实际中可能有负载均衡或选择策略 target_server servers[0] logging.info(fAttempting to connect to: {target_server[name]}) await connect_and_use_mcp_server(target_server) if __name__ __main__: asyncio.run(main())这个 Agent 演示了资源消费者的典型行为查询向 ARD 注册中心查询特定类型mcp-server和类别search的资源。选择从返回的列表中选择一个资源这里简单选择第一个。连接与调用根据资源信息这里需要扩展为实际的连接逻辑示例中简化了建立与 MCP 服务器的连接并调用其工具。4.4 运行与验证启动 ARD 注册中心在一个终端运行python ard_registry.py。启动 MCP 服务器在另一个终端运行python mcp_search_server.py。观察日志确认它成功注册。运行智能 Agent在第三个终端运行python smart_agent.py。观察日志你会看到它成功发现了注册的服务器并打印出模拟的搜索结果。这个原型虽然简单但它完整地演示了 ARD 的核心流程注册 - 发现 - 使用。你可以在此基础上扩展例如实现更复杂的资源筛选、负载均衡、安全认证在注册和发现时加入 Token或者将存储后端改为 Redis/PostgreSQL 以实现持久化和高可用。5. ARD 在真实场景中的挑战与进阶设计将 ARD 应用于生产环境或复杂项目时你会遇到一系列原型中未曾考虑的挑战。以下是几个关键问题及其解决思路。5.1 资源描述的标准化与语义发现原型中我们只用了简单的category和tags。但在真实世界里一个“搜索”服务可能有很多种网页搜索、学术搜索、内部文档搜索。一个“翻译”Agent 可能支持中英、中日互译。如何让消费者精确找到所需资源挑战简单的关键词匹配如category: “search”不够精确容易导致误匹配或找不到。解决方案富语义描述采用类似 OpenAPI Schema、Protocol Buffers.proto文件或自定义的 JSON Schema 来描述资源的接口工具的函数签名、输入输出格式。注册时将这些模式Schema的哈希或链接一同上传。能力本体Capability Ontology定义一个共享的、结构化的能力分类词汇表。例如使用一个共享的 URI 来表示https://capabilities.example/search#web和https://capabilities.example/translation#zh-en。Agent 在注册和查询时都引用这个标准词汇表。向量化搜索将资源的功能描述文本如工具的描述文档通过嵌入模型Embedding转换为向量存入向量数据库。消费者可以用自然语言如“找一个能总结长文章的助手”进行查询通过向量相似度匹配到最相关的资源。这为实现“语义发现”提供了可能。5.2 网络拓扑与连接建立我们的原型假设所有组件都在一个平坦的网络中且endpoint字段可以直接连接。这在容器化、多云或边缘计算环境中是不现实的。挑战网络隔离MCP 服务器运行在 Kubernetes Pod 内其 IP 对集群外不可达。动态地址Agent 是移动设备IP 地址经常变化。连接协议多样有的是 Stdio需要启动子进程有的是 TCP/HTTP/WebSocket。解决方案端点抽象与解析器endpoint字段不应是简单的地址字符串而应是一个抽象描述如{“type”: “stdio”, “command”: [“python”, “search_server.py”]}或{“type”: “k8s-service”, “service-name”: “tavily-mcp”, “port”: 8000}。ARD 客户端 SDK 需要根据type调用相应的“连接器”来建立实际连接。边车Sidecar模式为每个需要被发现的资源部署一个轻量级的“边车”代理。资源本身只与边车通信如通过本地 Unix Socket边车负责向 ARD 注册、维持心跳并对外提供统一的网络端点如负载均衡器后的地址。这是服务网格的常见模式。连接中继Relay对于无法直接寻址的资源如家庭网络中的设备可以通过一个公网可访问的“中继服务器”进行桥接。资源连接到中继并向 ARD 注册中继分配的通道 ID。消费者也通过中继与资源通信。5.3 安全、认证与授权开放式的发现意味着安全风险。不能允许任意 Agent 注册或调用任何服务。挑战如何防止恶意节点注册虚假服务如何控制哪些消费者可以发现和调用哪些资源解决方案双向 TLS 认证在注册和心跳阶段要求资源提供者出示由私有 CA 签发的客户端证书。ARD 注册中心只接受可信证书的注册。命名空间与多租户在 ARD 中引入“命名空间”或“项目”概念。资源注册在特定命名空间下。消费者只能发现其所属命名空间或经授权访问的命名空间内的资源。这类似于 Kubernetes 的 Namespace。能力声明与策略引擎资源注册时不仅声明能力还可能声明其所需的调用权限或敏感度标签如requires-auth: true,>
返回列表