Skip to content

Day 1 — Gateway 与消息协议

读完这章你能获得什么:理解 AI Agent 框架的通信基座——统一消息格式、事件发布订阅、消息路由分发、WebSocket 服务端,并能自己启动一个 Gateway 完成收发测试。


一、前情提要

这是我们的第一章,也是整个框架的地基。在做任何"智能"的事情之前,我们先要解决一个最基本的问题:

消息怎么进来?怎么出去?格式是什么?谁来处理?

如果你从 需求分析架构设计 读过来,你知道我们要造的是一个"私人助理"。但助理再聪明,如果连用户说的话都听不到,那也白搭。

本章要解决的就是"听"和"说"的基础设施。


二、生活类比:快递分拣中心

想象你开了一家快递公司的分拣中心。每天有成千上万的包裹从各地送来,你需要:

  1. 统一包装标签:不管是从哪个网点收来的包裹,都要贴上统一格式的标签(寄件人、收件人、类型、时间戳)—— 这就是 GatewayMessage
  2. 标签类型分类:标签上标注了"普通件"、"加急件"、"退货件"等类型 —— 这就是 MessageType
  3. 自动分拣机:根据标签类型,把包裹分到不同的处理流水线 —— 这就是 MessageRouter
  4. 广播系统:包裹到了通知仓库备货、新快递员上岗通知考勤系统——这些"旁路通知"不影响分拣主流程 —— 这就是 EventBus
  5. 分拣中心大楼:把上面所有东西装在一起,对外提供收发服务 —— 这就是 GatewayServer
现实世界                          代码世界
─────────                        ─────────
快递包裹                     →    消息(用户输入、Agent 回复等)
统一格式标签                 →    GatewayMessage(Pydantic 模型)
标签上的类型                 →    MessageType(text/command/event/error)
自动分拣机                   →    MessageRouter(按类型分发到处理器)
广播系统                     →    EventBus(发布-订阅,旁路通知)
分拣中心大楼                 →    GatewayServer(WebSocket 服务端)

三、本章要解决的问题

问题解决方案为什么需要
消息格式不统一,后续处理很麻烦定义 GatewayMessage 统一格式不管消息从哪来,内部都长一样
不同类型的消息需要不同处理方式MessageRouter 按类型分发文本消息和命令消息,处理逻辑完全不同
连接状态变化(上线、下线)需要通知其他组件EventBus 发布-订阅解耦:Gateway 不需要知道谁关心这些事件
需要一个网络服务来收发消息GatewayServer(WebSocket)客户端和服务端之间的实时通信通道

四、核心概念详解

4.1 GatewayMessage —— 统一的"快递标签"

所有在系统内部流转的消息,都被封装成同一种格式。就像快递公司不管你寄什么东西,包裹外面都要贴一张格式统一的面单。

消息有哪些字段?

字段类型类比说明
msg_idstr快递单号UUID 自动生成,全局唯一,用于追踪和日志
msg_typeMessageType快递类型(普通/加急/退货)决定消息被路由到哪个处理器
sourcestr寄件人消息从谁那来的(用户 ID、渠道名等)
targetstr收件人消息要发给谁(agent、某个服务等)
payloadAny包裹内容实际的业务数据(文本内容、事件数据等)
timestampdatetime寄件时间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 的区别是什么?

对比MessageRouterEventBus
模式请求-响应(一对一)发布-订阅(一对多)
返回值可以返回响应消息无返回值
用途处理业务消息(用户输入→回复)旁路通知(日志、监控、统计)
类比分拣机把包裹送到对应流水线广播系统通知各部门

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 适配器 将解决这个问题。