1. 模块概述
MessageBus(nanobot/bus/queue.py)是 nanobot 的通信中枢,通过两个 asyncio.Queue 实现聊天渠道与 Agent 核心的完全解耦。
📍 源码:
nanobot/nanobot/bus/queue.py#L1-L45
1.1 设计理念
MessageBus 的设计遵循 发布-订阅 模式的简化变体:
- 渠道 只负责将外部消息转换为
InboundMessage并发布到入站队列 - Agent 核心 只负责从入站队列消费消息,处理后发布
OutboundMessage到出站队列 - 渠道管理器 从出站队列消费消息并路由到对应渠道
这种解耦使得:
- 新增渠道不需要修改 Agent 核心代码
- Agent 核心可以独立测试,不需要真实渠道
- 消息处理失败不影响渠道连接
1.2 在系统中的位置
graph LR
subgraph "生产者"
TG[Telegram Channel]
DC[Discord Channel]
WS[WebSocket Channel]
CLI[CLI Direct]
end
subgraph "MessageBus"
IQ[Inbound Queue
asyncio.Queue]
OQ[Outbound Queue
asyncio.Queue]
end
subgraph "消费者"
AL[AgentLoop]
CM[ChannelManager]
end
TG -->|publish_inbound| IQ
DC -->|publish_inbound| IQ
WS -->|publish_inbound| IQ
CLI -->|process_direct| AL
IQ -->|consume_inbound| AL
AL -->|publish_outbound| OQ
OQ -->|consume_outbound| CM
CM -->|send| TG
CM -->|send| DC
CM -->|send| WS
2. 命名体系
MessageBus 的 API 命名遵循清晰的 发布/消费 对称模式:
| 操作 | 入站 | 出站 |
|---|---|---|
| 发布 | publish_inbound(msg) |
publish_outbound(msg) |
| 消费 | consume_inbound() → InboundMessage |
consume_outbound() → OutboundMessage |
| 查询 | inbound_size |
outbound_size |
命名规律:所有方法名采用 动词_方向 格式,语义自解释。
3. API Signatures
graph LR
subgraph "MessageBus"
PI[publish_inbound]
CI[consume_inbound]
PO[publish_outbound]
CO[consume_outbound]
IS[inbound_size]
OS[outbound_size]
end
Channel[Channel] --> PI
PI --> IQ[Inbound Queue]
IQ --> CI
CI --> Agent[AgentLoop]
Agent --> PO
PO --> OQ[Outbound Queue]
OQ --> CO
CO --> Manager[ChannelManager]
1 | class MessageBus: |
消息数据结构
1 |
|
4. 数据结构深度解析
4.1 MessageBus
4.a 结构体存在的理由
MessageBus 是整个系统的解耦点。没有它,渠道和 Agent 核心将直接耦合——每个渠道需要知道如何调用 Agent,Agent 需要知道如何向每个渠道发送回复。MessageBus 将 N×M 的耦合关系简化为 N+M:每个组件只需要知道如何与总线交互。
4.b 结构定义
1 | class MessageBus: |
4.c 字段三层分析表
| 字段 | 设计动机 | 反事实 | 替代方案 |
|---|---|---|---|
inbound |
解耦消息生产者和消费者,支持多生产者单消费者 | 渠道需要直接调用 Agent 方法,无法排队和背压 | 可以用回调注册,但队列提供自然的背压和排序 |
outbound |
解耦 Agent 和渠道分发,支持多消费者模式 | Agent 需要知道所有活跃渠道并逐一发送 | 可以用 pub/sub 中间件(Redis),但增加外部依赖 |
4.d 消息流转时序
sequenceDiagram
participant U as 用户
participant CH as Channel
participant MB as MessageBus
participant AL as AgentLoop
participant AR as AgentRunner
participant CM as ChannelManager
U->>CH: 发送消息
CH->>MB: publish_inbound(msg)
MB-->>AL: consume_inbound() 返回 msg
AL->>AR: _run_agent_loop(messages)
AR-->>AL: AgentRunResult
AL->>MB: publish_outbound(response)
MB-->>CM: consume_outbound() 返回 response
CM->>CH: send(response)
CH->>U: 显示回复
5. 函数逐行精讲
5.1 publish_inbound()
1 | async def publish_inbound(self, msg: InboundMessage) -> None: |
5.2 consume_inbound()
1 | async def consume_inbound(self) -> InboundMessage: |
6. 设计决策分析
6.1 为什么使用无界队列?
决策:asyncio.Queue() 不设置 maxsize,创建无界队列。
权衡:
- ✅ 不会因队列满而阻塞生产者(渠道),确保消息不丢失
- ✅ 简化错误处理——不需要处理
QueueFull异常 - ❌ 如果消费者(Agent)处理速度长期低于生产速度,内存可能无限增长
- ❌ 无法提供天然的背压信号
缓解措施:AgentLoop 使用 _concurrency_gate(默认 3)限制并发处理数,间接触发入站队列的背压。
6.2 为什么 Inbound 和 Outbound 是独立的两个队列?
决策:使用两个独立的 asyncio.Queue 分别处理入站和出站消息。
权衡:
- ✅ 入站和出站完全解耦:Agent 处理慢不影响出站消息发送
- ✅ 两个方向可以独立监控(
inbound_size/outbound_size) - ✅ 类型安全:
Queue[InboundMessage]和Queue[OutboundMessage]是不同的泛型类型 - ❌ 两个队列的状态不关联,无法实现请求-响应配对追踪
7. 学习检查点
📝 本章小结
- MessageBus 是两个 asyncio.Queue 的封装,设计极其简洁(仅 45 行代码)
- 发布/消费命名对称:
publish_inbound↔consume_inbound,publish_outbound↔consume_outbound - 无界队列设计:简化错误处理,配合并发门控实现隐式背压
- 解耦价值:将 N 渠道 × M 核心的耦合简化为 N+M
🤔 思考题
如果要将 MessageBus 改为有界队列以提供显式背压,
publish_inbound需要如何处理QueueFull异常?参考答案
有界队列满时,
put()会阻塞或(使用put_nowait()时)抛出QueueFull。渠道的_handle_message需要捕获此异常并决定策略:丢弃消息(记录警告日志)、等待重试(使用put()阻塞版本)、或返回错误给用户。最简单的实现是让publish_inbound使用await self.inbound.put(msg)(阻塞等待),这样背压会自然传递到渠道的消息接收循环。为什么
AgentLoop.run()使用wait_for(consume_inbound(), timeout=1.0)而不是直接await consume_inbound()?参考答案
1 秒超时(
nanobot/nanobot/agent/loop.py#L877)允许主循环在没有新消息时执行定期维护任务:auto_compact.check_expired()检查并压缩过期会话。如果使用无限阻塞的await,这些维护任务只能在有新消息到达时才能执行,在空闲时段过期会话永远不会被清理。