Environment 消息总线、Message 路由与 Team
30 秒导读: 单个角色的「观察-思考-行动」循环(01 章)只解决了「一个 agent 怎么自转」。多个 agent 要协作,就缺一样东西:它们之间怎么通信、消息怎么被正确投递给该收的人。这一章讲 MetaGPT 的答案——
Environment是一根消息总线:角色把消息「发布」到环境,环境按一本地址簿把它投递到目标角色的私有信箱;Team再在总线之上套一层预算 + 轮次,把整台机器驱动起来。
1. 这一章解决什么问题(零基础也能懂)
前两章的主角是单个角色:Role 是自转的原子(01 章),Action 是它手里的工作单元(02 章)。
但「多智能体」的关键从来不是「有很多角色」,而是角色之间怎么说话。设想产品经理写完了 PRD,他怎么把这份 PRD「递」给架构师、而不是递给测试工程师?架构师又怎么知道「这条消息是冲我来的、我该动手了」?
这就是本章要解决的两个问题:
- 投递问题: 一条消息发出去,谁该收到?怎么保证精确送达、不误投、不漏投?
- 驱动问题: 谁来推动这群角色一轮一轮地转下去,什么时候该停(活干完了 / 钱花光了)?
MetaGPT 的设计有一个关键的解耦思想(来自 RFC 113/116,见 §5):
消息只负责说清「发给谁」,完全不关心「对方在哪、怎么送过去」。 「怎么送达」是传输框架(即
Environment)的职责。
这正是 publish_message 源码注释里写的原话(metagpt/environment/base_env.py:175-183):路由信息只指定收件人,"without concern for where the message recipient is located"。
一句话直觉:把 Environment 想成一个公司内部的邮件系统。 你写邮件只填收件人名字(send_to),不用管对方坐哪、用什么网线——邮件系统(环境)拿着一本通讯录(member_addrs)负责把信塞进对方的收件箱(每个角色的私有 msg_buffer)。
2. 顶层全景(它大概怎么转)
三个主角
| 部件 | 白话职责 | 在哪个文件 |
|---|---|---|
Message | 一封信。带路由三元组 cause_by / sent_from / send_to,说明「因何而发、谁发的、发给谁」 | metagpt/schema.py:232 |
Environment | 邮件系统 / 消息总线。持有地址簿 member_addrs,负责分发与并发驱动所有角色 | metagpt/environment/base_env.py:124 |
Team | 公司老板。雇人(hire)、给预算(invest)、按轮次运转(run),并能存档/恢复整个项目 | metagpt/team.py:32 |
一条消息怎么从发出到送达
下图是本章的主线。怎么读:从上到下是一条消息的旅程,左边是「谁在做」,右边是「用哪个符号」。
Team.run_project(idea) team.py:102 run_project
│ 发布一封「用户需求」信 Message(content=idea)
▼
Environment.publish_message(msg) base_env.py:175
│
│ 遍历地址簿 member_addrs: {角色 → 它的地址集合}
│ 对每个角色问一句: is_send_to(msg, 它的地址)? common.py:423
│ ├─ msg.send_to 含 "<all>" → 命中(广播)
│ └─ 角色地址 ∩ msg.send_to ≠ ∅ → 命中(定向)
▼
role.put_message(msg) → 塞进该角色私有信箱 role.py:448 / msg_buffer
┊ (只有命中的角色才收到)
┊
▼ 下一轮 Environment.run() 时
role._observe(): 从信箱 pop_all,再按 cause_by/send_to 过滤出「我关心的」 role.py:399
│
▼
该角色进入「观察→思考→行动」循环(见 01 章),产出新 Message,再次 publish……
谁在推着轮子转
Environment 只管「投递」。真正让所有角色一轮一轮跑起来的,是两层驱动:
Environment.run()—— 一轮之内,把所有非空闲角色用asyncio.gather并发跑一遍(base_env.py:197)。Team.run(n_round)—— 外层大循环:每轮先查「全空闲了吗 / 钱够吗」,再调env.run(),直到轮数耗尽或提前停(team.py:122)。
3. 核心机制(逐个拆解)
3.1 Message:路由三元组 + check_* 归一化
它要解决的小问题: 一封信得能回答三个问题——「因 什么而发?谁发的?发给谁?」。这三个字段就是 MetaGPT 的路由三元组。
| 字段 | 常量名 | 含义 | 谁来读它 |
|---|---|---|---|
cause_by | MESSAGE_ROUTE_CAUSE_BY | 因何而发:触发这封信的那个 Action 的类全名 | 角色的 watch 订阅:「我盯着某类 Action 的产物」 |
sent_from | MESSAGE_ROUTE_FROM | 谁发的:发信角色的类全名 | 调试 / 溯源 |
send_to | MESSAGE_ROUTE_TO | 发给谁:一个收件人地址集合 set[str] | Environment 分发时匹配 |
这三个常量定义在 metagpt/const.py:77-79。注意 send_to 默认是 {MESSAGE_ROUTE_TO_ALL},即 "<all>"——不指定收件人就是广播(schema.py:241)。
巧妙处:字段一律「归一化成字符串」。 你在业务里可以把 cause_by 写成一个 Action 类、把 send_to 写成一个角色对象或名字,但它们进入 Message 时会被一组 check_* 校验器统一转成字符串 / 字符串集合,这样分发时只做纯字符串比较,简单可靠。
# 真实源码,schema.py:266-279 —— 三个字段的归一化校验器(节选)
@field_validator("cause_by", mode="before")
def check_cause_by(cls, cause_by: Any) -> str:
# 不填就默认「用户需求」;任何 Action 类都被 any_to_str 转成类全名字符串
return any_to_str(cause_by if cause_by else import_class("UserRequirement", ...))
@field_validator("send_to", mode="before")
def check_send_to(cls, send_to: Any) -> set:
# 单个/列表/集合都被转成「字符串集合」,不填则回落到 {<all>}
return any_to_str_set(send_to if send_to else {MESSAGE_ROUTE_TO_ALL})
any_to_str(common.py:395)的规则很朴素:字符串原样返回,类 / 对象则取类全名(如 metagpt.actions.write_prd.WritePRD)。所以「订阅某个 Action」和「发给某个 Action 类名」能对得上。
除三元组外还有几个关键字段:id(check_id 用 uuid4 自动补,schema.py:244)、role(system/user/assistant)、instruct_content(结构化产物,承载 02 章 的 ActionNode 输出)、metadata。
__setattr__ 的小陷阱(schema.py:307): 即使你事后赋值 msg.send_to = SomeRole,它也会被 __setattr__ 拦下来转成字符串集合——路由字段永远是干净的字符串,不会因为「后来手动改了一下」而破功。
便捷子类(schema.py:419-454): UserMessage / SystemMessage / AIMessage 只是预设了 role 的 Message,方便对接 OpenAI 消息格式。其中 AIMessage 多了 with_agent(name) / agent 属性——把「这条 AI 回复出自哪个角色」记进 metadata,供 05 章 的 TeamLeader 调度识别来源。
3.2 两个路由「魔法地址」:<all> 与 <self>
const.py:81-83 定义了三个特殊地址:
| 常量 | 值 | 作用 |
|---|---|---|
MESSAGE_ROUTE_TO_ALL | "<all>" | 广播:任何角色都算收件人 |
MESSAGE_ROUTE_TO_NONE | "<none>" | 谁都不发 |
MESSAGE_ROUTE_TO_SELF | "<self>" | 发给我自己 |
<all> 在 is_send_to 里被直接短路成「命中」(见 §3.3)。而 <self> 是在角色侧被翻译掉的——Role.publish_message(role.py:433-435)在发信前,把 send_to 里的 "<self>" 替换成发信角色自己的类全名:
# 真实源码,role.py:433-437 —— <self> 在离开角色时被就地翻译
if MESSAGE_ROUTE_TO_SELF in msg.send_to:
msg.send_to.add(any_to_str(self)) # 加上「我自己」的真实地址
msg.send_to.remove(MESSAGE_ROUTE_TO_SELF) # 摘掉占位符
if not msg.sent_from or msg.sent_from == MESSAGE_ROUTE_TO_SELF:
msg.sent_from = any_to_str(self) # 顺手补上「谁发的」
同一段还有个短路优化:如果这封信的收件人全是自己,角色直接 put_message 塞进自己信箱、根本不惊动环境(role.py:438-440)。只有要发给「别人」时,才真正 self.rc.env.publish_message(msg) 上总线。
3.3 Environment.publish_message:地址簿 + is_send_to 分发
它要解决的小问题: 一封信到了总线,怎么在一堆角色里挑出该收的那几个?
思路: 环境维护一本地址簿 member_addrs: Dict[BaseRole, Set](base_env.py:133)——每个角色对应它的一组地址。分发时逐个角色问 is_send_to:这封信的 send_to 和你的地址有交集吗?有就投。
# 真实源码,base_env.py:184-195 —— 分发的核心就这几行
def publish_message(self, message, peekable=True) -> bool:
found = False
for role, addrs in self.member_addrs.items(): # 遍历地址簿
if is_send_to(message, addrs): # 交集判定
role.put_message(message) # 命中 → 投进它的私有信箱
found = True
if not found:
logger.warning(f"Message no recipients: {message.dump()}")
self.history.add(message) # 全量留档(for debug)
return True
匹配逻辑 is_send_to(common.py:423-431)只有两条规则,先广播后定向:
def is_send_to(message, addresses: set):
if MESSAGE_ROUTE_TO_ALL in message.send_to: # ① 广播:含 <all> 直接命中
return True
for i in addresses: # ② 定向:地址与 send_to 有交集
if i in message.send_to:
return True
return False
两个值得记住的行为:
- 没有收件人不会报错,只
warning,并且消息照样进history——history是一份全量调试留档(base_env.py:134注释# For debug),不参与路由。 put_message只是把信push进角色的私有接收缓冲rc.msg_buffer(role.py:448-452),并不立即处理。真正消费发生在下一轮该角色_observe时(见 §3.6 的说明与 01 章)。
3.4 地址簿从哪来:add_roles / set_addresses
它要解决的小问题: member_addrs 这本通讯录是谁、什么时候填的?
答案藏在「雇人」的链路里。Team.hire → Environment.add_roles(base_env.py:164):
# 真实源码,base_env.py:164-173 —— 加入一批角色
def add_roles(self, roles):
for role in roles:
self.roles[role.name] = role # 先登记进 roles 名册
for role in roles:
role.context = self.context # 共享全局 Context
role.set_env(self) # 关键:反向把「环境」交给角色
role.set_env(env) 会回调 env.set_addresses(self, self.addresses)(role.py:308-313),把这个角色 → 它的地址集合写进 member_addrs(base_env.py:240-242)。那么一个角色的默认地址是什么?看 Role.check_addresses(role.py:211-213):
默认地址 =
{角色类全名, 角色名},例如产品经理是{"metagpt.roles.product_manager.ProductManager", "Alice"}。
这解释了整套路由为什么能「对上」: 你只要把 send_to 写成对方的类名或名字,就能命中它的地址集合。地址簿 = 角色注册时用「类名 + 名字」自动登记的双键索引。
set_addresses 也可被角色主动调用(role.py:293-300)来改自己的订阅地址,改完立刻同步回环境,地址簿始终最新。
3.5 Environment.run:并发跑非空闲角色 + is_idle + archive
它要解决的小问题: 一轮里,怎么让该干活的角色一起动、干完的别空转?
# 真实源码,base_env.py:197-211 —— 一轮驱动
async def run(self, k=1):
for _ in range(k):
futures = []
for role in self.roles.values():
if role.is_idle: # 空闲的跳过,不浪费
continue
futures.append(role.run()) # 收集协程
if futures:
await asyncio.gather(*futures) # 并发跑这一批
三个要点:
- 并发而非串行。 一轮里所有非空闲角色用
asyncio.gather同时推进;彼此靠上一轮投递到信箱的消息解耦,不需要显式排队。 is_idle决定谁参与、以及全局是否停。 环境的is_idle(base_env.py:228-234)是「所有角色都空闲」才为真;而单个角色的is_idle(role.py:557-559)= 没有待处理新消息(news)、没有 todo、私有信箱msg_buffer也空。信箱一旦被投进新信,该角色就不再空闲,下一轮自动被拉起来。archive:收尾存档(base_env.py:244-247)。 项目跑完时,若context里带project_path,就对生成的代码仓做一次GitRepository.archive()——把「一行需求跑成一个软件仓库」(04 章)的产物提交归档。
3.6 一封信从「投进信箱」到「被消费」
投递(§3.3)只把消息塞进 rc.msg_buffer,没有过滤。真正的「这封信我要不要理」发生在角色的 _observe(role.py:399-419):
# 真实源码,role.py:406-412(节选)—— 从信箱取信,再按订阅过滤
news = self.rc.msg_buffer.pop_all() # 清空信箱,取出所有新信
self.rc.news = [
n for n in news
if (n.cause_by in self.rc.watch # 我 watch 的 Action 触发的?
or self.name in n.send_to) # 或点名发给我的?
and n not in old_messages # 且没处 理过(去重)
]
于是 MetaGPT 的路由其实是两级筛:
- 环境级(粗筛,§3.3):
send_to ∩ 地址—— 决定「信进不进你信箱」。 - 角色级(细筛,这里):
cause_by ∈ watch或name ∈ send_to—— 决定「进了信箱的信,你这轮理不理」。
cause_by 的价值在此凸显:角色不必知道谁会产出它要的东西,只需 watch 某类 Action——基于「因何而发」订阅,而非基于「谁发的」。这是 SOP 流水线(04 章)能自然串起来的底层原因。
3.7 Team:预算、轮次、存档
Team(team.py:32)是最外层的老板,给消息总线套上经营约束。
雇人 hire(team.py:83): 就是转调 env.add_roles,把角色接进环境地址簿(§3.4)。
给钱 invest(team.py:92-96): 把投资额写进 CostManager.max_budget。运行时每轮 _check_balance(team.py:98-100)检查累计花费是否超预算,超了就抛 NoMoneyException——LLM 调用是要花钱的,这是硬性刹车。
# 真实源码,team.py:98-100 —— 预算刹车
def _check_balance(self):
if self.cost_manager.total_cost >= self.cost_manager.max_budget:
raise NoMoneyException(self.cost_manager.total_cost, f"Insufficient funds: ...")
发起项目 run_project(team.py:102-107): 把一行需求包成一封普通 Message(content=idea)(默认广播、cause_by 归一为 UserRequirement)发上总线——这就是整条 SOP 的第一封信。
主循环 run(team.py:122-138):
# 真实源码,team.py:128-137(节选)
while n_round > 0:
if self.env.is_idle: # 全员空闲 → 活干完了,提前收工
break
n_round -= 1
self._check_balance() # 钱不够 → 抛异常止损
await self.env.run() # 驱动一轮(§3.5)
self.env.archive(auto_archive) # 收尾归档
return self.env.history # 返回全量消息留档
停机有两个条件:轮数 n_round 耗尽,或 env.is_idle(没有任何消息在流动了)。@serialize_decorator(team.py:122)会在异常时自动把现场序列化,便于恢复。
存档 / 恢复(team.py:59-81): serialize 把整个 Team(含 env.context)写成 team.json;deserialize 反向重建——整台多智能体机器的状态可持久化、可断点续跑。