1. 介绍README.md# 智能售后中枢ANP A2A MCP演示项目 ## 场景 用户发起售后换货请求时系统同时使用三种协议 | 协议 | 职责 | 本项目对应 | |------|------|------------| | **ANP** | 服务发现与负载均衡 | anp_routing.py 选择负载最低的客服节点 | | **A2A** | 多专家 Agent 协作 | a2a_experts.py 技术 / 政策 / 仓储 | | **MCP** | 访问业务系统工具 | mcp_biz_server.py 订单 / 物流 / 工单 | ## 架构 text 用户请求 ↓ ANPcs_node_* 集群选路最低负载 ↓ 编排器main.py ├─ MCPquery_order / query_logistics / open_ticket └─ A2Atech.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(nametech_diagnosis, description技术诊断专家) policy A2AServer(nameafter_sales_policy, description售后政策专家) warehouse A2AServer(namewarehouse_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_asciiFalse) 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(dayswarranty_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_asciiFalse) 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: fEX{order_id[-4:]}8888, message: 已生成退换运单预计 3 日内寄出新品, } return json.dumps(result, ensure_asciiFalse) 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(host127.0.0.1, portport) threads [ threading.Thread(target_run, args(tech, TECH_PORT), daemonTrue), threading.Thread(target_run, args(policy, POLICY_PORT), daemonTrue), threading.Thread(target_run, args(warehouse, WAREHOUSE_PORT), daemonTrue), ] 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(backgroundTrue) clients get_expert_clients() demo clients[tech].execute_skill( diagnose, json.dumps({product: 耳机, issue: 右耳无声坏了}, ensure_asciiFalse), ) 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( discoverydiscovery, service_idfcs_node_{i}, service_namef客服节点{i}, service_typecustomer_service, capabilities[after_sales, order_support, exchange], endpointfhttp://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_typecustomer_service) if not nodes: return None return min(nodes, keylambda 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( nameafter-sales-biz, description售后业务系统 MCP 服务订单/物流/工单, ) def query_order(order_id: str) - str: 根据订单号查询订单和客户信息。 return json.dumps(get_order(order_id), ensure_asciiFalse, indent2) def query_logistics(order_id: str) - str: 根据订单号查询物流信息。 return json.dumps(get_logistics(order_id), ensure_asciiFalse, indent2) def open_ticket(order_id: str, issue: str, action: str 换货) - str: 创建售后工单。 return json.dumps( create_ticket(order_id, issue, action), ensure_asciiFalse, indent2, ) def list_all_tickets() - str: 列出当前已创建的工单。 return json.dumps(list_tickets(), ensure_asciiFalse, indent2) 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, encodingutf-8) as f: return json.load(f) except (json.JSONDecodeError, OSError): return [] def _save_tickets(tickets: list) - None: with open(_TICKET_FILE, w, encodingutf-8) as f: json.dump(tickets, f, ensure_asciiFalse, indent2) 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: fTK{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( nameafter_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(rORD\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]}, fregion{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_asciiFalse) ) 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_asciiFalse) ) 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_asciiFalse) ) 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_count5) 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_typecustomer_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_count5) 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 (idcs_node_3, load0.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 (idcs_node_2, load0.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 (idcs_node_3, load0.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: load0.36 - 客服节点1: load0.57 - 客服节点2: load0.15 - 客服节点3: load0.17 - 客服节点4: load0.44