尧图网站建设 尧图网络
  • 首页
  • 关于我们
  • 服务项目
  • 案例展示
  • 建站流程
  • 资讯中心
  • 联系我们
首页/资讯中心/详情

agent-MCP-A2A-ANP销售售后项目

agent-MCP-A2A-ANP销售售后项目
📅 发布时间:2026/7/23 9:51:15

1. 介绍

README.md

# 智能售后中枢(ANP + A2A + MCP)演示项目 ## 场景 用户发起售后换货请求时,系统同时使用三种协议: | 协议 | 职责 | 本项目对应 | |------|------|------------| | **ANP** | 服务发现与负载均衡 | `anp_routing.py` 选择负载最低的客服节点 | | **A2A** | 多专家 Agent 协作 | `a2a_experts.py` 技术 / 政策 / 仓储 | | **MCP** | 访问业务系统工具 | `mcp_biz_server.py` 订单 / 物流 / 工单 | ## 架构 ```text 用户请求 ↓ ANP:cs_node_* 集群选路(最低负载) ↓ 编排器(main.py) ├─ MCP:query_order / query_logistics / open_ticket └─ A2A:tech.diagnose → policy.review → warehouse.arrange ↓ 汇总售后结论 ``` ## 目录 ```text after_sales_hub/ ├── mock_data.py # 模拟订单/客户/物流数据 ├── mcp_biz_server.py # MCP 业务工具服务 ├── a2a_experts.py # A2A 专家服务 ├── anp_routing.py # ANP 注册与选路 ├── main.py # 完整编排入口 └── README.md ``` ## 运行 在 `agent` conda 环境中: ```bash cd /data/after_sales_hub # 自动跑 3 个演示案例 python main.py # 交互模式 python main.py interactive ``` 可选:单独启动 MCP / A2A(一般不必,main 会自动拉起): ```bash python mcp_biz_server.py python a2a_experts.py ``` ## 演示订单 | 订单号 | 商品 | 说明 | |--------|------|------| | `ORD20260301` | 无线降噪耳机 Pro | 在保,适合换货演示 | | `ORD20260302` | 智能手表 S2 | 可能过保,政策可能拒绝 | ## 说明 1. 本项目是教学演示:订单数据为内存模拟,不连真实数据库。 2. 编排器用代码串联三协议,避免 LLM 漏传工具参数导致演示失败。 3. A2A 专家服务监听 `7100/7101/7102`;若端口占用请先释放。

2. a2a_experts

# @File : a2a_experts.py """A2A 专业客服智能体:技术诊断 / 售后政策 / 仓储换货。 职责(A2A):多个专家 Agent 点对点协作处理复杂售后问题。 """ from __future__ import annotations import json import re import threading import time from datetime import datetime, timedelta from hello_agents.protocols import A2AServer, A2AClient TECH_PORT = 7100 POLICY_PORT = 7101 WAREHOUSE_PORT = 7102 def _parse_json_payload(text: str) -> dict: """从技能输入中尽量解析 JSON。""" text = text.strip() # 兼容 "answer {...}" / "diagnose {...}" 前缀 for prefix in ("answer ", "diagnose ", "review ", "arrange "): if text.lower().startswith(prefix): text = text[len(prefix) :].strip() break try: return json.loads(text) except json.JSONDecodeError: match = re.search(r"\{.*\}", text, re.DOTALL) if match: try: return json.loads(match.group(0)) except json.JSONDecodeError: pass return {"raw": text} def create_experts() -> tuple[A2AServer, A2AServer, A2AServer]: tech = A2AServer(name="tech_diagnosis", description="技术诊断专家") policy = A2AServer(name="after_sales_policy", description="售后政策专家") warehouse = A2AServer(name="warehouse_exchange", description="仓储换货专家") @tech.skill("diagnose") def diagnose(text: str) -> str: data = _parse_json_payload(text) issue = str(data.get("issue", data.get("raw", ""))) product = data.get("product", "未知商品") quality_keywords = ["坏了", "损坏", "无声", "断连", "无法开机", "质量"] is_quality = any(k in issue for k in quality_keywords) result = { "expert": "tech_diagnosis", "product": product, "issue": issue, "is_quality_issue": is_quality, "confidence": 0.9 if is_quality else 0.6, "suggestion": "建议走质量换货" if is_quality else "建议先远程排障", } return json.dumps(result, ensure_ascii=False) @policy.skill("review") def review(text: str) -> str: data = _parse_json_payload(text) order = data.get("order", {}) diagnosis = data.get("diagnosis", {}) purchase_date = order.get("purchase_date") warranty_days = int(order.get("warranty_days", 0)) in_warranty = False if purchase_date: start = datetime.strptime(purchase_date, "%Y-%m-%d") in_warranty = datetime.now() <= start + timedelta(days=warranty_days) approved = bool(diagnosis.get("is_quality_issue")) and in_warranty result = { "expert": "after_sales_policy", "in_warranty": in_warranty, "approved": approved, "action": "换货" if approved else "拒绝换货", "reason": ( "在保且判定为质量问题,同意换货" if approved else "不在保或非质量问题,需人工复核" ), } return json.dumps(result, ensure_ascii=False) @warehouse.skill("arrange") def arrange(text: str) -> str: data = _parse_json_payload(text) order_id = data.get("order_id", "UNKNOWN") approved = bool(data.get("approved", False)) logistics = data.get("logistics", {}) if not approved: result = { "expert": "warehouse_exchange", "arranged": False, "message": "政策未通过,不安排换货", } else: result = { "expert": "warehouse_exchange", "arranged": True, "order_id": order_id, "return_carrier": logistics.get("carrier", "默认快递"), "exchange_tracking": f"EX{order_id[-4:]}8888", "message": "已生成退换运单,预计 3 日内寄出新品", } return json.dumps(result, ensure_ascii=False) return tech, policy, warehouse def start_experts(background: bool = True) -> None: tech, policy, warehouse = create_experts() def _run(server: A2AServer, port: int) -> None: server.run(host="127.0.0.1", port=port) threads = [ threading.Thread(target=_run, args=(tech, TECH_PORT), daemon=True), threading.Thread(target=_run, args=(policy, POLICY_PORT), daemon=True), threading.Thread(target=_run, args=(warehouse, WAREHOUSE_PORT), daemon=True), ] for t in threads: t.start() time.sleep(2) print("✅ A2A 专家服务已启动: tech:7100 / policy:7101 / warehouse:7102") def get_expert_clients() -> dict[str, A2AClient]: return { "tech": A2AClient("http://127.0.0.1:7100"), "policy": A2AClient("http://127.0.0.1:7101"), "warehouse": A2AClient("http://127.0.0.1:7102"), } if __name__ == "__main__": start_experts(background=True) clients = get_expert_clients() demo = clients["tech"].execute_skill( "diagnose", json.dumps({"product": "耳机", "issue": "右耳无声坏了"}, ensure_ascii=False), ) print(demo) try: while True: time.sleep(1) except KeyboardInterrupt: print("\nA2A 服务已停止")

3. anp_routing

# @File : anp_routing.py """ANP 客服集群:服务注册、发现与负载均衡。 职责(ANP):在大规模并发下发现可用客服节点并选择负载最低者。 """ from __future__ import annotations import random from hello_agents.protocols import ANPDiscovery, register_service def build_cs_cluster(node_count: int = 5) -> ANPDiscovery: discovery = ANPDiscovery() for i in range(node_count): register_service( discovery=discovery, service_id=f"cs_node_{i}", service_name=f"客服节点{i}", service_type="customer_service", capabilities=["after_sales", "order_support", "exchange"], endpoint=f"http://cs-node-{i}:8000", metadata={ "load": round(random.uniform(0.1, 0.6), 2), "region": random.choice(["华北", "华东", "华南"]), "max_concurrency": random.choice([50, 80, 100]), }, ) return discovery def pick_best_node(discovery: ANPDiscovery): """选择负载最低的客服节点(Least Load)。""" nodes = discovery.discover_services(service_type="customer_service") if not nodes: return None return min(nodes, key=lambda s: s.metadata.get("load", 1.0)) def bump_load(node, delta: float = 0.08) -> None: """模拟节点接单后负载上升。""" current = float(node.metadata.get("load", 0.0)) node.metadata["load"] = round(min(current + delta, 0.99), 2) def release_load(node, delta: float = 0.05) -> None: """模拟处理完成后负载下降。""" current = float(node.metadata.get("load", 0.0)) node.metadata["load"] = round(max(current - delta, 0.05), 2)

4. mcp_biz_server

# @File : mcp_biz_server.py #!/usr/bin/env python3 """MCP 业务工具服务:订单 / 物流 / 工单。 职责(MCP):让智能体通过标准工具接口访问业务系统。 """ import json import os import sys sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from hello_agents.protocols import MCPServer from mock_data import get_order, get_logistics, create_ticket, list_tickets biz_server = MCPServer( name="after-sales-biz", description="售后业务系统 MCP 服务(订单/物流/工单)", ) def query_order(order_id: str) -> str: """根据订单号查询订单和客户信息。""" return json.dumps(get_order(order_id), ensure_ascii=False, indent=2) def query_logistics(order_id: str) -> str: """根据订单号查询物流信息。""" return json.dumps(get_logistics(order_id), ensure_ascii=False, indent=2) def open_ticket(order_id: str, issue: str, action: str = "换货") -> str: """创建售后工单。""" return json.dumps( create_ticket(order_id, issue, action), ensure_ascii=False, indent=2, ) def list_all_tickets() -> str: """列出当前已创建的工单。""" return json.dumps(list_tickets(), ensure_ascii=False, indent=2) biz_server.add_tool(query_order) biz_server.add_tool(query_logistics) biz_server.add_tool(open_ticket) biz_server.add_tool(list_all_tickets) if __name__ == "__main__": biz_server.run()

5. mock_data

# @File : mock_data.py """模拟客户、订单、物流数据(演示用,不连真实数据库)。""" from __future__ import annotations import json import os from copy import deepcopy from threading import Lock CUSTOMERS = { "U1001": {"name": "张三", "phone": "138****0001", "level": "金牌会员"}, "U1002": {"name": "李四", "phone": "139****0002", "level": "普通会员"}, } ORDERS = { "ORD20260301": { "order_id": "ORD20260301", "user_id": "U1001", "product": "无线降噪耳机 Pro", "price": 899.0, "purchase_date": "2026-07-10", "warranty_days": 365, "status": "已完成", }, "ORD20260302": { "order_id": "ORD20260302", "user_id": "U1002", "product": "智能手表 S2", "price": 1299.0, "purchase_date": "2025-12-01", "warranty_days": 365, "status": "已完成", }, } LOGISTICS = { "ORD20260301": { "order_id": "ORD20260301", "carrier": "顺丰速运", "tracking_no": "SF1234567890", "status": "已签收", "signed_at": "2026-07-12", }, "ORD20260302": { "order_id": "ORD20260302", "carrier": "京东物流", "tracking_no": "JD9876543210", "status": "已签收", "signed_at": "2025-12-05", }, } _TICKET_FILE = os.path.join(os.path.dirname(__file__), ".tickets.json") _LOCK = Lock() def _load_tickets() -> list: if not os.path.exists(_TICKET_FILE): return [] try: with open(_TICKET_FILE, "r", encoding="utf-8") as f: return json.load(f) except (json.JSONDecodeError, OSError): return [] def _save_tickets(tickets: list) -> None: with open(_TICKET_FILE, "w", encoding="utf-8") as f: json.dump(tickets, f, ensure_ascii=False, indent=2) def get_order(order_id: str) -> dict: order = ORDERS.get(order_id) if not order: return {"error": f"订单不存在: {order_id}"} result = deepcopy(order) result["customer"] = deepcopy(CUSTOMERS.get(order["user_id"], {})) return result def get_logistics(order_id: str) -> dict: info = LOGISTICS.get(order_id) if not info: return {"error": f"未找到物流信息: {order_id}"} return deepcopy(info) def create_ticket(order_id: str, issue: str, action: str) -> dict: with _LOCK: tickets = _load_tickets() ticket = { "ticket_id": f"TK{len(tickets) + 1:04d}", "order_id": order_id, "issue": issue, "action": action, "status": "已创建", } tickets.append(ticket) _save_tickets(tickets) return deepcopy(ticket) def list_tickets() -> list: with _LOCK: return deepcopy(_load_tickets())

6. main

# @File : main.py #!/usr/bin/env python3 """智能售后中枢:同时使用 ANP + A2A + MCP 的完整演示。 请求链路: 用户请求 → ANP 选择负载最低的客服节点 → MCP 查询订单 / 物流,并创建工单 → A2A 技术 / 政策 / 仓储专家协作 → 汇总回复 """ from __future__ import annotations import json import os import re import sys BASE_DIR = os.path.dirname(os.path.abspath(__file__)) sys.path.insert(0, BASE_DIR) from hello_agents.tools import MCPTool from a2a_experts import get_expert_clients, start_experts from anp_routing import build_cs_cluster, bump_load, pick_best_node, release_load MCP_SERVER = os.path.join(BASE_DIR, "mcp_biz_server.py") def create_mcp_tool() -> MCPTool: return MCPTool( name="after_sales_biz", description="售后业务系统:订单、物流、工单", server_command=["python", MCP_SERVER], ) def mcp_call(mcp: MCPTool, tool_name: str, arguments: dict) -> dict: raw = mcp.run( { "action": "call_tool", "tool_name": tool_name, "arguments": arguments, } ) # MCPTool 返回值可能带前缀说明,尽量抽出 JSON try: return json.loads(raw) except json.JSONDecodeError: match = re.search(r"\{.*\}|\[.*\]", raw, re.DOTALL) if match: try: return json.loads(match.group(0)) except json.JSONDecodeError: pass return {"raw": raw} def extract_order_id(text: str) -> str | None: match = re.search(r"ORD\d+", text.upper()) return match.group(0) if match else None def handle_after_sales_request( user_query: str, discovery, mcp: MCPTool, experts: dict, ) -> str: print("\n" + "=" * 60) print(f"📥 用户请求: {user_query}") print("=" * 60) # ---------- 1) ANP:选节点 ---------- node = pick_best_node(discovery) if not node: return "❌ ANP 未发现可用客服节点" bump_load(node) print( f"[ANP] 路由到 {node.service_name} " f"(id={node.service_id}, load={node.metadata['load']}, " f"region={node.metadata.get('region')})" ) try: order_id = extract_order_id(user_query) or "ORD20260301" issue = user_query # ---------- 2) MCP:查订单 / 物流 ---------- print(f"[MCP] query_order({order_id})") order = mcp_call(mcp, "query_order", {"order_id": order_id}) if "error" in order: return f"❌ 订单查询失败: {order['error']}" print(f"[MCP] query_logistics({order_id})") logistics = mcp_call(mcp, "query_logistics", {"order_id": order_id}) # ---------- 3) A2A:专家协作 ---------- tech_payload = { "product": order.get("product"), "issue": issue, } print("[A2A] tech.diagnose ...") tech_resp = experts["tech"].execute_skill( "diagnose", json.dumps(tech_payload, ensure_ascii=False) ) diagnosis = json.loads(tech_resp.get("result", "{}")) policy_payload = {"order": order, "diagnosis": diagnosis} print("[A2A] policy.review ...") policy_resp = experts["policy"].execute_skill( "review", json.dumps(policy_payload, ensure_ascii=False) ) policy = json.loads(policy_resp.get("result", "{}")) warehouse_payload = { "order_id": order_id, "approved": policy.get("approved", False), "logistics": logistics, } print("[A2A] warehouse.arrange ...") wh_resp = experts["warehouse"].execute_skill( "arrange", json.dumps(warehouse_payload, ensure_ascii=False) ) warehouse = json.loads(wh_resp.get("result", "{}")) # ---------- 4) MCP:创建工单 ---------- action = policy.get("action", "人工复核") print(f"[MCP] open_ticket(action={action})") ticket = mcp_call( mcp, "open_ticket", { "order_id": order_id, "issue": issue, "action": action, }, ) # ---------- 5) 汇总 ---------- customer = order.get("customer", {}) summary = f""" ✅ 售后处理完成(节点: {node.service_name}) 【客户】{customer.get('name', '未知')}({customer.get('level', '')}) 【订单】{order_id} / {order.get('product')} / ¥{order.get('price')} 【物流】{logistics.get('carrier', '-')} {logistics.get('tracking_no', '-')}({logistics.get('status', '-')}) 【技术诊断】质量问题={diagnosis.get('is_quality_issue')};建议={diagnosis.get('suggestion')} 【售后政策】在保={policy.get('in_warranty')};结论={policy.get('action')};原因={policy.get('reason')} 【仓储安排】{warehouse.get('message')} 【工单】{ticket.get('ticket_id', ticket)} """.strip() print(summary) return summary finally: release_load(node) print( f"[ANP] 节点 {node.service_name} 处理结束," f"当前负载={node.metadata['load']}" ) def demo(): print("🚀 启动智能售后中枢(ANP + A2A + MCP)") discovery = build_cs_cluster(node_count=5) print(f"✅ ANP 已注册 {len(discovery.list_all_services())} 个客服节点") start_experts() experts = get_expert_clients() mcp = create_mcp_tool() # 先验证 MCP 工具列表 tools = mcp.run({"action": "list_tools"}) print("\n[MCP] 可用工具预览:") print(tools[:500] if isinstance(tools, str) else tools) cases = [ "我的订单 ORD20260301 耳机坏了,右耳无声,想申请换货", "订单 ORD20260302 手表无法开机,质量有问题,请求换货", "ORD20260301 想了解一下物流到哪了,顺便问问能不能换货,耳机断连了", ] for query in cases: handle_after_sales_request(query, discovery, mcp, experts) print("\n" + "=" * 60) print("📊 最终各客服节点负载:") for svc in discovery.discover_services(service_type="customer_service"): print(f" - {svc.service_name}: load={svc.metadata['load']}") print("=" * 60) def interactive(): print("🚀 启动智能售后中枢(交互模式)") print("示例: 我的订单 ORD20260301 耳机坏了,想换货") print("输入 quit 退出\n") discovery = build_cs_cluster(node_count=5) start_experts() experts = get_expert_clients() mcp = create_mcp_tool() while True: query = input("你: ").strip() if not query: continue if query.lower() in {"quit", "exit", "q"}: break reply = handle_after_sales_request(query, discovery, mcp, experts) print("\n助手:\n" + reply + "\n") if __name__ == "__main__": if len(sys.argv) > 1 and sys.argv[1] == "interactive": interactive() else: demo()

7. 测试

============================================================ 📥 用户请求: 我的订单 ORD20260301 耳机坏了,右耳无声,想申请换货 ============================================================ [ANP] 路由到 客服节点3 (id=cs_node_3, load=0.19, region=华东) [MCP] query_order(ORD20260301) ✅ 连接成功! INFO:mcp.server.lowlevel.server:Processing request of type CallToolRequest INFO:mcp.server.lowlevel.server:Processing request of type ListToolsRequest 🔌 连接已断开 ✅ 售后处理完成(节点: 客服节点3) 【客户】张三(金牌会员) 【订单】ORD20260301 / 无线降噪耳机 Pro / ¥899.0 【物流】顺丰速运 SF1234567890(已签收) 【技术诊断】质量问题=True;建议=建议走质量换货 【售后政策】在保=True;结论=换货;原因=在保且判定为质量问题,同意换货 【仓储安排】已生成退换运单,预计 3 日内寄出新品 【工单】TK0001 [ANP] 节点 客服节点3 处理结束,当前负载=0.14
============================================================ 📥 用户请求: 订单 ORD20260302 手表无法开机,质量有问题,请求换货 ============================================================ [ANP] 路由到 客服节点2 (id=cs_node_2, load=0.2, region=华南) [MCP] query_order(ORD20260302) 连接成功! INFO:mcp.server.lowlevel.server:Processing request of type CallToolRequest INFO:mcp.server.lowlevel.server:Processing request of type ListToolsRequest 🔌 连接已断开 ✅ 售后处理完成(节点: 客服节点2) 【客户】李四(普通会员) 【订单】ORD20260302 / 智能手表 S2 / ¥1299.0 【物流】京东物流 JD9876543210(已签收) 【技术诊断】质量问题=True;建议=建议走质量换货 【售后政策】在保=True;结论=换货;原因=在保且判定为质量问题,同意换货 【仓储安排】已生成退换运单,预计 3 日内寄出新品 【工单】TK0002 [ANP] 节点 客服节点2 处理结束,当前负载=0.15
============================================================ 📥 用户请求: ORD20260301 想了解一下物流到哪了,顺便问问能不能换货,耳机断连了 ============================================================ [ANP] 路由到 客服节点3 (id=cs_node_3, load=0.22, region=华东) [MCP] query_order(ORD20260301) 📝 使用 Stdio 传输 (命令): python 🔗 连接到 MCP 服务器... ✅ 连接成功! INFO:mcp.server.lowlevel.server:Processing request of type CallToolRequest INFO:mcp.server.lowlevel.server:Processing request of type ListToolsRequest 🔌 连接已断开 ✅ 售后处理完成(节点: 客服节点3) 【客户】张三(金牌会员) 【订单】ORD20260301 / 无线降噪耳机 Pro / ¥899.0 【物流】顺丰速运 SF1234567890(已签收) 【技术诊断】质量问题=True;建议=建议走质量换货 【售后政策】在保=True;结论=换货;原因=在保且判定为质量问题,同意换货 【仓储安排】已生成退换运单,预计 3 日内寄出新品 【工单】TK0003 [ANP] 节点 客服节点3 处理结束,当前负载=0.17 ============================================================ 📊 最终各客服节点负载: - 客服节点0: load=0.36 - 客服节点1: load=0.57 - 客服节点2: load=0.15 - 客服节点3: load=0.17 - 客服节点4: load=0.44 ============================================================

相关新闻

  • AI+隐患排查|金汤令双AI引擎构建智慧隐患治理体系
  • KVM虚拟化技术详解:从部署到优化实践
  • 2026 怀化黄金回收新规解读,本地十年老店永兴回收测评指南 - 黄金珠宝

最新新闻

  • CNN-LSTM混合模型在风电功率预测中的应用与优化
  • TI BQ41Z50状态寄存器深度解析:从OperationStatus到GaugingStatus的实战调试指南
  • 把人肉流程抽成脚本:重复操作识别与可回放工具化
  • 虚拟桌面切换
  • Git版本控制器
  • TUSB系列8052芯片无JTAG调试:串口打印与Keil ISD51实战指南

日新闻

  • 亨得利盐城维修点在哪里?手表维修保养地址指南**公示(2026年7月最新) - 亨得利官方
  • 提升.NET API安全性:Boxed.AspNetCore.Swagger认证授权最佳实践
  • 帝舵佛山**网点地址更新:2026年7月售后热线电话与服务客户指南 - 帝舵中国官方服务中心

周新闻

  • SaaS软件行业GEO实践:AI搜索时代的品牌可见性与获客新路径
  • 什么是PCTFE?医药高端包装的“防潮王牌“材料
  • 【JVM调优实战】16-可视化利器-JConsole-VisualVM-JMC

月新闻

  • 2026年6月公司网站搭建最新热门渠道测评:四大低成本/零代码平台对比+避坑
  • 【Linux】Linux arm 编译QT程序,出现expected “}“报错
  • 【MATLAB例程】四基站二维AOA定位与距离辅助增强对比仿真。基于角度观测和测距修正的固定目标平面定位精度分析

关于尧图

  • 公司简介
  • 团队介绍
  • 企业文化
  • 荣誉资质

服务项目

  • 定制开发
  • 电商建站
  • UI 设计
  • 运维服务

快速链接

  • 案例展示
  • 建站流程
  • 常见问题
  • 资讯中心

联系方式

  • 📍北京市朝阳区互联网产业园 A 座 10 层
  • 📞400-888-8888
  • ✉️contact@rkmt.cn
  • 🕐周一至周日 9:00-21:00

© 2024 北京尧图网络科技有限公司 版权所有 | 京 ICP 备 XXXXXXXX 号