乐于分享
好东西不私藏

openclaw源码解读:Gateway+消息总线

openclaw源码解读:Gateway+消息总线

导读: 上一篇我们聊了龙虾(OpenClaw)的整体架构,这一篇小编直接把 nanobot 的源码摊开给你看——重点拆解 Gateway 网关层和消息总线(Bus),用大白话讲清楚每一层在干嘛。零基础友好,附完整源码解读。


一、麻雀虽小,五脏俱全

上次聊完龙虾的整体架构,后台好几条留言问我同一个问题:"道理我都懂,但代码到底长什么样?"

行,今天就拿 OpenClaw 生态里最精简的项目——nanobot——来做一次"开膛手术"。

nanobot 只有 5000 来行代码。猜猜它支持多少聊天平台?

19 个。

Telegram、Discord、Slack、微信、飞书、钉钉、QQ、WhatsApp、Email……你叫得上名字的 IM,它基本都接了。

5000 行。19 个平台。

我第一反应跟你一样:"这代码得烂成什么样?"

然后点开源码——哎,还真不烂。不光不烂,甚至称得上优雅。秘密就在今天要讲的两个模块:Gateway(网关) 和 Bus(消息总线)


二、Gateway 网关:一套接口打天下

2.1 核心思路——"插件化"接入

你想啊,Telegram 用 HTTP 长轮询,Discord 走 WebSocket,Email 走 IMAP/SMTP,微信还得扫码登录……每个平台的协议完全不一样。

硬写 if-else?写到第 5 个平台你就想离职了。

nanobot 的解法很直接:搞一个抽象基类 BaseChannel,所有平台继承它,实现三个方法就完事。

说白了就是策略模式 + 插件架构,但用得恰到好处。

2.2 BaseChannel——网关的"宪法"

直接上源码,你感受一下这个基类有多克制:

classBaseChannel(ABC):"""所有聊天平台的抽象基类"""    name: str = "base"# 平台标识,如 "telegram"    display_name: str = "Base"# 展示名称    send_progress: bool = True# 是否发送进度提示    show_reasoning: bool = True# 是否展示推理过程def__init__(self, config: Any, bus: MessageBus):self.config = configself.bus = bus          # 消息总线,后面重点讲self._running = False    @abstractmethodasyncdefstart(self) -> None:"""启动频道,开始监听消息(长期运行的异步任务)"""pass    @abstractmethodasyncdefstop(self) -> None:"""停止频道,清理资源"""pass    @abstractmethodasyncdefsend(self, msg: OutboundMessage) -> None:"""通过该平台发送消息"""pass

就三个抽象方法:start()stop()send()

翻译成人话:

  • • start():连接平台,开始听消息
  • • stop():断开连接,收工回家
  • • send():把 AI 的回复发出去

任何一个平台,只要实现这三个方法,就能接入 nanobot。就这么简单。

2.3 权限校验——一个方法搞定

BaseChannel 里还有个隐藏的"保镖"——_handle_message() 方法。所有平台收到用户消息后,都统一调用它:

asyncdef_handle_message(    self, sender_id, chat_id, content, media=None,    metadata=None, session_key=None, is_dm=False) -> None:"""统一入口:校验权限 → 构建消息 → 发布到总线"""# 第一关:权限检查ifnotself.is_allowed(sender_id):if is_dm:# 私聊的陌生人?给他一个配对码            code = generate_code(self.name, str(sender_id))awaitself.send(OutboundMessage(                channel=self.name,                chat_id=str(chat_id),                content=format_pairing_reply(code),            ))return# 没权限,到此为止# 第二关:构建标准消息格式    msg = InboundMessage(        channel=self.name,        sender_id=str(sender_id),        chat_id=str(chat_id),        content=content,        media=media or [],        metadata=meta,    )# 第三关:扔到消息总线上awaitself.bus.publish_inbound(msg)

看懂了吗?逻辑就三步:检查权限 → 封装成统一格式 → 扔到消息总线

不管你是从 Telegram 来的还是从 Discord 来的,到了总线上,长得一模一样。各个平台的脏活累活——协议转换、媒体下载、格式适配——全在各自的 Channel 里消化干净,出来的就是标准的 InboundMessage

2.4 is_allowed()——权限三级检查

defis_allowed(self, sender_id: str) -> bool:"""权限检查:通配符 > 白名单 > 配对码 > 拒绝"""    allow_list = getattr(self.config, "allow_from", [])if"*"in allow_list:      # 通配符:所有人都能用returnTrueifstr(sender_id) in allow_list:  # 在白名单里returnTrueif is_approved(self.name, str(sender_id)):  # 用配对码验证过的returnTruereturnFalse# 没过任何关卡?拒绝

三层递进,很好理解:

  1. 1. 配置了 *?全放行,谁来都能用
  2. 2. sender_id 在白名单里?放行
  3. 3. 之前用配对码验证过?放行
  4. 4. 以上都没过?那不好意思了

2.5 自动发现机制——零配置接入

你可能好奇:框架怎么知道有哪些 Channel 可用?

答案是 registry.py——基于 pkgutil 的自动发现

defdiscover_channel_names() -> list[str]:"""零导入扫描:只读文件名,不导入任何 SDK"""import nanobot.channels as pkgreturn [        namefor _, name, ispkg in pkgutil.iter_modules(pkg.__path__)if name notin _INTERNAL andnot ispkg    ]

这个函数只扫描 channels/ 目录下有哪些 .py 文件,不 import 任何一个

为什么?因为每个 Channel 都依赖不同的三方 SDK——Telegram 用 python-telegram-bot,Discord 用 discord.py,Slack 用 slack-sdk……如果全部 import 一遍,启动速度会慢到你怀疑人生。

所以它的策略是:先看菜单(文件名),用户点了哪个才去做(import)

defload_channel_class(module_name: str) -> type[BaseChannel]:"""按需导入:只加载用户启用的 Channel"""    mod = importlib.import_module(f"nanobot.channels.{module_name}")for attr indir(mod):        obj = getattr(mod, attr)ifisinstance(obj, typeandissubclass(obj, BaseChannel) and obj isnot BaseChannel:return objraise ImportError(f"找不到 BaseChannel 子类")

这才是真正的懒加载——不是偷懒,是聪明。你只用了 Telegram 和微信,凭什么要等 Discord SDK 加载完?

2.6 ChannelManager——串起来的那个人

BaseChannel 定义接口,Registry 负责发现,但谁来把这些 Channel 真正跑起来、管起来?

ChannelManager。它干的活不少:

  1. 1. 发现并初始化所有已启用的 Channel
  2. 2. 并发启动所有 Channel
  3. 3. 分发出站消息到对应的 Channel
  4. 4. 失败重试(指数退避:1s → 2s → 4s)
  5. 5. 流式消息合并(多个 delta 合成一次发送,减少 API 调用)
  6. 6. 去重(防止回声循环)
classChannelManager:asyncdefstart_all(self) -> None:"""启动所有频道和出站分发器"""# 1. 启动出站消息分发循环self._dispatch_task = asyncio.create_task(self._dispatch_outbound())# 2. 并发启动所有 Channel        tasks = []for name, channel inself.channels.items():            tasks.append(asyncio.create_task(self._start_channel(name, channel)))await asyncio.gather(*tasks, return_exceptions=True)

注意那个 _dispatch_outbound()——它是一个无限循环的后台任务,专门从消息总线的出站队列里取消息,然后路由到对应的 Channel。


三、消息总线(Bus)——解耦的艺术

3.1 为什么需要消息总线?

想象一下没有消息总线的世界:

Telegram Channel 收到消息 → 直接调用 Agent 处理 → Agent 处理完直接调用 Telegram Channel 发送

看起来也行?但问题来了——如果你同时接了 19 个平台呢?

每个 Channel 都要知道 Agent 在哪、怎么调,Agent 也要知道所有 Channel 的发送接口。这叫紧耦合——改一个地方,到处都得跟着改。

消息总线的做法是:所有人只跟总线打交道,互相不认识。

就像小区的快递柜——你不用跟快递员面对面交接,往柜子里一放,对方自己来取。

3.2 MessageBus——45 行,没了

我说整个消息总线的核心就 45 行代码,你可能不信。看:

classMessageBus:"""异步消息总线:解耦聊天频道和 Agent 核心"""def__init__(self):# 两条队列,一进一出self.inbound: asyncio.Queue[InboundMessage] = asyncio.Queue()self.outbound: asyncio.Queue[OutboundMessage] = asyncio.Queue()asyncdefpublish_inbound(self, msg: InboundMessage) -> None:"""频道往里塞消息(用户发的)"""awaitself.inbound.put(msg)asyncdefconsume_inbound(self) -> InboundMessage:"""Agent 从里面取消息(阻塞等待)"""returnawaitself.inbound.get()asyncdefpublish_outbound(self, msg: OutboundMessage) -> None:"""Agent 往里塞回复"""awaitself.outbound.put(msg)asyncdefconsume_outbound(self) -> OutboundMessage:"""频道从里面取回复(阻塞等待)"""returnawaitself.outbound.get()

两个队列,四个方法。没了。

你可能觉得:"就这?" 对,就这。但就是这么短的代码,把 Channel 和 Agent 彻底解耦了:

  • • Channel 只管往 inbound 队列里塞,不关心谁来消费
  • • Agent 只管从 inbound 取、往 outbound 塞,不关心谁来发送
  • • ChannelManager 从 outbound 取出来,路由到对应 Channel

三方各做各的事,互不干扰。

3.3 消息格式——两个 dataclass 定义一切

@dataclassclassInboundMessage:"""用户发来的消息"""    channel: str# 来自哪个平台    sender_id: str# 谁发的    chat_id: str# 在哪个会话    content: str# 说了什么    timestamp: datetime = field(default_factory=datetime.now)    media: list[str] = field(default_factory=list)    # 附件    metadata: dict[strAny] = field(default_factory=dict)    @propertydefsession_key(self) -> str:"""会话标识:平台名:会话ID"""returnself.session_key_override orf"{self.channel}:{self.chat_id}"@dataclassclassOutboundMessage:"""发给用户的回复"""    channel: str# 发到哪个平台    chat_id: str# 发到哪个会话    content: str# 回复内容    reply_to: str | None = None    media: list[str] = field(default_factory=list)    metadata: dict[strAny] = field(default_factory=dict)    buttons: list[list[str]] = field(default_factory=list)  # 交互按钮

设计思路就一句话:把所有平台的消息"翻译"成同一种语言。

Telegram 的 Update 对象也好、Discord 的 Message 也好、邮件的 MIME 格式也好——进了总线,统统变成 InboundMessage

这里有个细节值得说说:那个 metadata 字段。它是个开放的字典,平台特有的东西可以随便往里塞——Discord 的 guild_id、Slack 的 thread_ts,都放这里。不污染核心字段,但需要的时候随时能取。这种"预留扩展口但不过度设计"的做法,我个人挺喜欢的。

3.4 RuntimeEventBus——给前端开的"后门"

除了主总线,还有个"副总线"——RuntimeEventBus

这个不走队列,走的是观察者模式,就是经典的 Pub/Sub:

classRuntimeEventBus:"""进程内事件总线,用于运行时状态通知"""def__init__(self) -> None:self._handlers: list[_HandlerEntry] = []defsubscribe(self, handler, event_type=None):"""订阅事件,可选过滤类型"""        entry = (event_type, handler)self._handlers.append(entry)# 返回取消订阅的函数def_unsubscribe():self._handlers.remove(entry)return _unsubscribeasyncdefpublish(self, event) -> None:"""发布事件,按注册顺序通知所有匹配的订阅者"""for event_type, handler inlist(self._handlers):if event_type isnotNoneandnotisinstance(event, event_type):continue# 类型不匹配,跳过try:                result = handler(event)if inspect.isawaitable(result):await resultexcept Exception:                logger.exception("事件处理失败")

为什么需要两套总线?

因为它们的职责完全不同:

MessageBus(主总线)
RuntimeEventBus(副总线)
用途
用户消息收发
系统状态通知
模式
队列(一对一消费)
Pub/Sub(一对多广播)
消费者
Agent / ChannelManager
WebUI / 监控面板
典型事件
用户发了一条消息
Agent 开始处理 / 处理完成

RuntimeEventBus 上跑的都是些"边角料"——"Agent 开始思考了""这轮用了 2.3 秒""模型切换了"。WebUI 订阅这些事件来做实时状态展示,但就算你把 WebUI 关了,聊天功能一点不受影响。

说白了,主总线管"干活",副总线管"汇报进度"。各管各的,互不影响。


四、消息流转全景——从按下回车到收到回复

前面讲了一堆模块,现在把它们串起来。一条消息从用户手里发出去,到收到 AI 回复,中间经过了什么?

入站流程(用户 → AI):

入站流程图

出站流程(AI → 用户):

出站流程图

还有个骚操作——流式合并:

你用过 ChatGPT 吧?它是一个字一个字蹦出来的。nanobot 也支持这种效果,AI 每生成一小段就发一个带 _stream_delta 标记的消息。

但问题是——如果 AI 生成速度比网络发送速度快呢?队列里堆了 20 个 delta,你总不能调 20 次 API 吧?

所以 ChannelManager 干了一件聪明事:把队列里连续的 delta 合并成一次发送

def_coalesce_stream_deltas(self, first_msg):"""合并连续的流式片段,减少 API 调用次数"""    combined_content = first_msg.contentwhileTrue:try:            next_msg = self.bus.outbound.get_nowait()except asyncio.QueueEmpty:break# 同一个会话的 delta?合并内容if same_target and is_delta:            combined_content += next_msg.contentif is_end:  # 流结束标记breakelse:break# 不是同一个流了,停止合并return merged_message

这段逻辑不复杂,但解决了一个很实际的性能问题。没有这个合并,Telegram Bot API 的速率限制分分钟给你触发。


五、写在最后

说实话,读完这两个模块的源码,我最大的感受不是"哇好厉害",而是"原来可以这么简单"。

很多人(包括以前的我)做多平台接入,第一反应是搞一套超复杂的抽象层,定义几十个接口方法,考虑各种未来可能出现的场景。

结果呢?写了一堆用不上的代码,真正要接新平台的时候还是得改一堆东西。

nanobot 反过来——基类就仨方法,消息总线就俩队列,自动发现就扫一下文件名。每个模块都克制到了极点,但组合起来,19 个平台稳稳接住。

这大概就是所谓的"少即是多"吧。

如果你自己在做 AI Agent 框架、或者任何需要接多端的系统,这套"抽象基类 + 消息总线 + 懒加载"的三件套可以直接抄作业。


看完这篇,你更想看 Agent 核心的推理循环,还是 Session 会话管理的实现?评论区说一声,下一篇就拆它。

觉得有用就点个关注在看——下期见 👋