主题
Day 1 — Gateway 与消息协议
读完这章你能获得什么:理解 AI Agent 框架的通信基座——统一消息格式、事件发布订阅、消息路由分发、WebSocket 服务端,并能自己启动一个 Gateway 完成收发测试。
一、前情提要
这是我们的第一章,也是整个框架的地基。在做任何"智能"的事情之前,我们先要解决一个最基本的问题:
消息怎么进来?怎么出去?格式是什么?谁来处理?
如果你从 需求分析 和 架构设计 读过来,你知道我们要造的是一个"私人助理"。但助理再聪明,如果连用户说的话都听不到,那也白搭。
本章要解决的就是"听"和"说"的基础设施。
二、生活类比:快递分拣中心
想象你开了一家快递公司的分拣中心。每天有成千上万的包裹从各地送来,你需要:
- 统一包装标签:不管是从哪个网点收来的包裹,都要贴上统一格式的标签(寄件人、收件人、类型、时间戳)—— 这就是 GatewayMessage
- 标签类型分类:标签上标注了"普通件"、"加急件"、"退货件"等类型 —— 这就是 MessageType
- 自动分拣机:根据标签类型,把包裹分到不同的处理流水线 —— 这就是 MessageRouter
- 广播系统:包裹到了通知仓库备货、新快递员上岗通知考勤系统——这些"旁路通知"不影响分拣主流程 —— 这就是 EventBus
- 分拣中心大楼:把上面所有东西装在一起,对外提供收发服务 —— 这就是 GatewayServer
现实世界 代码世界
───────── ─────────
快递包裹 → 消息(用户输入、Agent 回复等)
统一格式标签 → GatewayMessage(Pydantic 模型)
标签上的类型 → MessageType(text/command/event/error)
自动分拣机 → MessageRouter(按类型分发到处理器)
广播系统 → EventBus(发布-订阅,旁路通知)
分拣中心大楼 → GatewayServer(WebSocket 服务端)三、本章要解决的问题
| 问题 | 解决方案 | 为什么需要 |
|---|---|---|
| 消息格式不统一,后续处理很麻烦 | 定义 GatewayMessage 统一格式 | 不管消息从哪来,内部都长一样 |
| 不同类型的消息需要不同处理方式 | MessageRouter 按类型分发 | 文本消息和命令消息,处理逻辑完全不同 |
| 连接状态变化(上线、下线)需要通知其他组件 | EventBus 发布-订阅 | 解耦:Gateway 不需要知道谁关心这些事件 |
| 需要一个网络服务来收发消息 | GatewayServer(WebSocket) | 客户端和服务端之间的实时通信通道 |
四、核心概念详解
4.1 GatewayMessage —— 统一的"快递标签"
所有在系统内部流转的消息,都被封装成同一种格式。就像快递公司不管你寄什么东西,包裹外面都要贴一张格式统一的面单。
消息有哪些字段?
| 字段 | 类型 | 类比 | 说明 |
|---|---|---|---|
msg_id | str | 快递单号 | UUID 自动生成,全局唯一,用于追踪和日志 |
msg_type | MessageType | 快递类型(普通/加急/退货) | 决定消息被路由到哪个处理器 |
source | str | 寄件人 | 消息从谁那来的(用户 ID、渠道名等) |
target | str | 收件人 | 消息要发给谁(agent、某个服务等) |
payload | Any | 包裹内容 | 实际的业务数据(文本内容、事件数据等) |
timestamp | datetime | 寄件时间 | UTC 时间戳,统一用 UTC 避免时区混乱 |
消息类型(MessageType)有哪些?
| 类型 | 值 | 用途 | 类比 |
|---|---|---|---|
| TEXT | "text" | 普通对话消息 | 普通快递 |
| COMMAND | "command" | 系统指令(如 /help、/reset) | 加急件 |
| EVENT | "event" | 内部事件通知 | 内部通知单 |
| ERROR | "error" | 错误信息 | 问题件 |
关键代码导读(miniclaw/gateway/protocol.py):
python
class MessageType(StrEnum):
"""用 StrEnum 而不是普通 Enum——这样枚举值直接就是字符串,
序列化为 JSON 时不需要任何转换。"""
TEXT = "text"
COMMAND = "command"
EVENT = "event"
ERROR = "error"
class GatewayMessage(BaseModel):
"""Pydantic 模型会自动帮我们做格式校验:
如果有人传了一个不存在的 msg_type,立刻报错,而不是让脏数据往下流。"""
msg_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
msg_type: MessageType = MessageType.TEXT
source: str = ""
target: str = ""
payload: Any = None
timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))还提供了快捷构造方法,省得每次都手填所有字段:
python
# 构造一条文本消息(最常用)
msg = GatewayMessage.text(content="你好", source="user_123", target="agent")
# 构造一条错误消息
err = GatewayMessage.error("消息处理失败", source="router")4.2 EventBus —— 广播系统
EventBus 是什么?
一句话:EventBus(事件总线)是一个中间人,让系统中的各个组件不需要互相认识,就能互相通信。
你可以把它想象成一个公告栏:任何人都可以在上面贴通知(发布事件),任何人也可以订阅某类通知(注册监听器)。贴通知的人不需要知道谁在看,看通知的人也不需要知道是谁贴的。
为什么需要 EventBus?—— 解耦
假设没有 EventBus,当用户连接时要做三件事:记日志、发欢迎消息、更新在线列表。代码会变成:
python
# 没有 EventBus:所有逻辑耦合在一起
# handle_connection 必须"认识"日志模块、消息模块、状态模块
# 改任何一个都要动这里
async def handle_connection(user):
await log_connection(user) # 日志模块
await send_welcome(user) # 消息模块
await update_online_list(user) # 状态模块有了 EventBus 之后:
python
# 各模块独立注册,互不知道对方的存在
bus.on("client_connected", log_connection)
bus.on("client_connected", send_welcome)
bus.on("client_connected", update_online_list)
# 触发时只需一行,不需要知道谁在听
await bus.emit("client_connected", user=user)每个模块只和 EventBus 打交道,互相之间完全不知道对方的存在。要加新功能?再 bus.on(...) 注册一个就行,不用改任何已有代码。
三个核心操作:
| 方法 | 作用 | 类比 |
|---|---|---|
on(event, handler) | 订阅:我关心某个事件,发生时请通知我 | 在公告栏前留下联系方式 |
emit(event, **kwargs) | 发布:某个事件发生了,通知所有关心它的人 | 在公告栏上贴通知 |
off(event, handler) | 取消订阅:我不再关心这个事件了 | 撤掉联系方式 |
关键设计——失败隔离:如果某个订阅者的处理函数出错了(比如门禁系统崩了),不会影响其他订阅者——考勤系统照常工作。这是因为 emit 内部捕获了每个处理器的异常并记录日志,而不是让异常冒泡出去。
4.3 MessageRouter —— 自动分拣机
为什么需要路由而不是一堆 if-else?
没有 Router 的代码:
python
# 丑陋的写法——每加一种消息类型就要改这里
if msg.msg_type == "text":
handle_text(msg)
elif msg.msg_type == "command":
handle_command(msg)
elif msg.msg_type == "event":
handle_event(msg)
# ...越来越长有 Router 的代码:
python
# 优雅的写法——注册一次,以后自动路由
router.register(MessageType.TEXT, handle_text)
router.register(MessageType.COMMAND, handle_command)
# 新增类型?只需要再 register 一行Router 和 EventBus 的区别是什么?
| 对比 | MessageRouter | EventBus |
|---|---|---|
| 模式 | 请求-响应(一对一) | 发布-订阅(一对多) |
| 返回值 | 可以返回响应消息 | 无返回值 |
| 用途 | 处理业务消息(用户输入→回复) | 旁路通知(日志、监控、统计) |
| 类比 | 分拣机把包裹送到对应流水线 | 广播系统通知各部门 |
4.4 GatewayServer —— 分拣中心大楼
GatewayServer 把上面所有组件组装在一起,对外提供 WebSocket 服务:
连接管理:GatewayServer 内部用一个 set 保存所有活跃连接——set 的添加和删除都是 O(1)(瞬间完成),适合频繁的连接/断开操作。
五、动手实验指南
5.1 运行示例
在项目根目录执行:
bash
python day1-gateway/example/main.py你会看到类似输出:
INFO Gateway 已启动: ws://127.0.0.1:xxxxx
INFO [事件] client_connected: ('127.0.0.1', 54321)
INFO [客户端] 收到: {"msg_type":"text","payload":"echo: hello miniOpenClaw",...}
INFO [事件] client_disconnected: ('127.0.0.1', 54321)
INFO Gateway 已停止5.2 改一改,看看会怎样
实验 1:修改 echo_handler,让它不只是简单回显,而是在内容前面加上当前时间。
实验 2:注册一个 COMMAND 类型的处理器——当用户发送 /status 时返回当前连接数。
实验 3:试试发送一个格式错误的 JSON 给 Gateway,观察它如何返回错误消息而不是崩溃。
5.3 运行测试
bash
pytest day1-gateway/gateway/ -v六、常见问题 FAQ
Q1:为什么 msg_id 用 UUID 而不是自增数字?
A:因为系统是分布式的——多个客户端同时连上来,如果用自增数字会冲突。UUID 是全局唯一的,不需要中心化协调。类比:身份证号 vs 排队号。排队号在不同窗口可能重复,身份证号全国唯一。
Q2:EventBus 和消息队列(RabbitMQ/Kafka)是一回事吗?
A:理念相似(都是发布-订阅),但我们的 EventBus 是进程内的——所有监听器在同一个 Python 进程里。消息队列是跨进程/跨机器的。教学阶段用进程内 EventBus 就够了,生产环境可以换成消息队列。
Q3:为什么用 WebSocket 而不是简单的 HTTP?
A:HTTP 是"你问我答"——客户端问了服务端才回。WebSocket 是"开着的电话"——双方随时说话。AI Agent 需要主动推送消息(比如流式输出每个字),HTTP 做不到。详见 技术调研。
七、知识拓展
- 与 Discord Bot 的对比:Discord Bot 的 Gateway 也是 WebSocket 长连接,消息格式也有 type/payload——思路完全一致,只是消息格式不同
- 与微服务网关(如 Kong/Nginx)的对比:微服务网关是 HTTP 层的路由分发,我们的 Gateway 是 WebSocket 层的——层次不同但理念相同
- 设计模式:EventBus 用的是"观察者模式",Router 用的是"策略模式"——两种经典的设计模式
下一章
你已经有了消息进出的基础设施。但目前用户只能通过 WebSocket 直连 Gateway——如果用户想通过命令行、网页、微信来聊天呢?
下一章 Day 2: Channel 适配器 将解决这个问题。