跳转至

Swarm nodes cn

Swarm 节点是 Jianmu 多 Agent 协作的行为树原语集合,它将 AgentRuntime 的底层消息路由、Agent 生命周期管理和邮箱同步机制封装为可组合的 py_trees 节点。开发者通过在行为树中编排 SpawnAgent、SendMessage、WaitMessage 三个核心节点,即可构建出从简单的 Coordinator-Worker 协作到液态拓扑(liquid-topology)自组织网络在内的各类多 Agent 通信模式。这三个节点分别对应多 Agent 系统中最基础的三个动作:创建新 Agent、发送消息和等待接收消息——它们共同构成了一套完整的 CSP(Communicating Sequential Processes)风格的 Agent 间通信协议。

架构定位:节点如何接入 Swarm 运行时

在深入各节点实现之前,有必要理解它们在整体架构中的位置。Swarm 节点并不直接处理消息的存储、路由或投递——这些职责由 AgentRuntime 承担。每个 AgentRuntime 实例内部维护一个 InMemoryMessageStore(消息持久化)、三套路由索引(_group_members 组播、_topic_subscribers 主题订阅、直接点对点)以及一个 AgentHost 注册表。节点通过 _SwarmRuntimeNode 基类提供的运行时解析机制,在 update_async 时动态获取与之关联的 AgentRuntime 引用,随后调用其 spawn()、send_message() 等公开方法完成实际操作。

下面的 Mermaid 图展示了三个核心节点与 AgentRuntime 内部组件之间的调用关系:

graph TB
    subgraph "行为树层 (Behavior Tree)"
        SA[SpawnAgent]
        SM[SendMessage]
        WM[WaitMessage]
    end

    subgraph "AgentRuntime"
        SP["spawn()"]
        SND["send_message()"]
        WK["wake()"]
        MSG[InMemoryMessageStore]
        HOST[AgentHost 注册表]
        MBX[mailbox 同步]
    end

    subgraph "路由层 (Routing)"
        RT[routing.send_message]
        ENQ[enqueue_hot_envelope]
        GRP[join_group / leave_group]
        TOP[subscribe_topic]
    end

    SA -->|"runtime.spawn(role, task)"| SP
    SP --> HOST
    SM -->|"runtime.send_message(...)"| SND
    SND --> RT
    RT --> MSG
    RT --> ENQ
    ENQ --> HOST
    ENQ --> WK
    WK --> MBX
    MBX -->|"写入 ctx.incoming_messages"| WM
    WM -->|"检查 ctx.incoming_messages"| MBX

关键点在于:SendMessage 将消息写入消息存储并通过 enqueue_hot_envelope 直接推送到目标 AgentHost 的 inbox 队列,随后 wake() 触发目标 Agent 的邮箱同步;WaitMessage 则通过检查 ctx.incoming_messages 来感知新消息的到达——这个属性正是由 mailbox 同步过程写入的。

_SwarmRuntimeNode:运行时解析基类

三个 Swarm 节点都直接或间接继承自 _SwarmRuntimeNode,这个内部基类继承 AsyncNode,提供了两个关键的辅助方法:

_resolve_runtime() 采用三级回退策略查找 AgentRuntime 实例:优先使用构造时传入的 runtime 属性;若为 None,从状态管理器中读取 "runtime" 键(支持全局回退);最后尝试从整个状态对象上通过 getattr(state, "runtime") 获取。解析后还会验证运行时是否具备 required_methods 中指定的方法——例如 SpawnAgent 要求 "spawn",SendMessage 要求 "send_message" 和 "emit"。

_read_state_value() 封装了带命名空间解析的状态读取逻辑,先通过 resolve_input_port 解析键名,然后在端口绑定的状态键上读取;若未找到且 allow_global_fallback=True,还会在全局状态中查找。

这种设计使得节点既可以在构造时显式绑定运行时(适用于静态流程),也可以依赖状态管理器在运行时动态获取(适用于 SwarmRole 等模板化角色场景)。

SpawnAgent:Agent 生成原语

SpawnAgent 是动态创建 Agent 的行为树节点。它调用 AgentRuntime.spawn() 完成实际的 Host 创建、行为树构建和状态初始化,并将新生成的 agent_id 写入状态以供后续节点使用。

参数体系

参数 类型 默认值 说明
role str \| None None 静态角色名,设置后忽略 role_key 查找
role_key str "role" 从状态中读取角色名的键
task str \| None None 静态任务描述,设置后忽略 task_key 查找
task_key str "task" 从状态中读取任务的键
parent_id str \| None None 静态父 Agent ID,设置后忽略 parent_id_key
parent_id_key str "agent_id" 从状态中读取父 Agent ID 的键
mode str "detached" 生成模式:"detached"(独立)或 "attached"(依附父 Agent)
constraints object \| None None 静态执行约束,覆盖 constraints_key 查找
constraints_key str "constraints" 从状态中读取约束的键
output_key str "spawned_agent_id" 新 Agent ID 写入状态的键
reuse_if_present bool False 若为 True,当 output_key 已有值时复用而不重新生成

执行逻辑

update_async() 首先解析运行时引用——若运行时不存在或不具备 spawn 方法,直接返回 FAILURE。接着检查 reuse_if_present 模式:如果启用了复用且 output_key 对应的状态槽中已存在非空字符串,则维持该值并返回 SUCCESS,实现幂等的 Agent 创建语义。

在正常生成路径中,节点依次解析 role(必填,缺失则 FAILURE)、task、parent_id 和 constraints,随后调用 runtime.spawn(role=..., task=..., parent_id=..., mode=..., constraints=...)。AgentRuntime.spawn() 内部执行以下步骤:

  1. 解析角色对象并生成 UUID 作为 agent_id
  2. 构建 AgentProfile 并调用角色的 build_tree(profile) 构建行为树根节点
  3. 通过 state_factory 创建 StateManager 并注入初始状态(agent_id、role、task)
  4. 组装 RunContext(包含 model_client、runtime、constraints、sandbox 等)
  5. 创建 AgentHost 并注册到 _hosts 字典
  6. 若 auto_start=True,立即启动该 Agent 的后台运行任务

生成的 agent_id 通过 write_port("spawned_agent_id", agent_id, state_key=self.output_key) 写入状态,后续节点通过读取该键获取目标 Agent 标识。

attached 模式与父 Agent 生命周期

当 mode="attached" 且 parent_id 非空时,子 Agent 会被记录到父 Agent 的 _attachments 集合中。这意味着当父 Agent 被 kill 时,所有 attached 子 Agent 也会被级联终止。"detached" 模式下的 Agent 则拥有独立的生命周期。

SendMessage:消息发送原语

SendMessage 是向其他 Agent、组或主题发送消息的行为树节点。它封装了 AgentRuntime.send_message() 的完整路由语义,支持三种目标寻址方式。

参数体系

参数 类型 默认值 说明
to_agent_id str \| None None 静态直接消息目标 Agent ID
to_agent_key str "to_agent_id" 从状态读取目标 Agent ID 的键
group_id str \| None None 静态组播目标组 ID
group_key str "group_id" 从状态读取组 ID 的键
topic str \| None None 静态主题名称
topic_key str "topic" 从状态读取主题的键
content str \| None None 静态消息正文,设置后忽略 content_key
content_key str "message" 从状态读取消息正文的键
sender_id str \| None None 静态发送者 ID,设置后忽略 sender_key
sender_key str "agent_id" 从状态读取发送者 ID 的键

路由语义

SendMessage 要求 to_agent_id、group_id、topic 三者恰好有一个非空——这是由底层 routing.send_message() 强制的约束:若零个或多个目标被同时指定,将抛出 ValueError。这种严格的路由排他性避免了消息被意外广播或路由歧义。

双路径发送策略

update_async() 实现了两条发送路径,体现了渐进增强的设计理念:

主路径(runtime.send_message):当解析到的运行时具备 send_message 方法时(AgentRuntime 实例具备),直接调用 await runtime.send_message(...),消息经过完整的存储持久化 → 信封入队 → 目标唤醒流程。

回退路径(事件发射):当运行时仅具备 emit 方法但不具备 send_message 时(例如某些精简运行时实现),节点自行构造 Message 对象并封装为 AgentEvent(event_type="message_sent", ...) 通过 runtime.emit() 发布。这条路径确保节点在最小运行时契约下仍可工作,但不会触发自动的消息存储和目标唤醒。

这种双路径策略使 SendMessage 保持了与 AgentRuntimeProtocol 中不同实现层级的兼容性。

消息投递的完整链路

当使用 AgentRuntime 的主路径时,一次消息发送的完整链路如下:

sequenceDiagram
    participant SM as SendMessage
    participant RT as AgentRuntime
    participant Routing as routing.send_message
    participant Store as InMemoryMessageStore
    participant Enq as enqueue_hot_envelope
    participant Host as 目标 AgentHost
    participant MBX as mailbox.sync_messages
    participant WM as WaitMessage (目标)

    SM->>RT: send_message(sender_id, content, to_agent_id)
    RT->>Routing: send_message(...)
    Routing->>Store: append(envelope) → seq
    Routing->>Enq: enqueue_hot_envelope(env)
    Enq->>Host: host.inbox.append(envelope)
    Enq->>Host: mark_dirty + wake
    Routing->>RT: emit("message_sent")
    RT->>MBX: schedule mailbox sync
    MBX->>Host: read inbox → set ctx.incoming_messages
    WM->>WM: 检查 ctx.incoming_messages → SUCCESS

WaitMessage:消息等待原语

WaitMessage 是所有 Swarm 节点中唯一继承同步 Node(而非 AsyncNode)的节点。它不执行异步 I/O,而是通过轮询 ctx.incoming_messages 属性来判断是否有新消息到达。

设计原理

WaitMessage 的同步特性源于 mailbox 系统的前置工作模式:在每次 Agent 的 tick 周期开始前,调度器(scheduler.step_agent)会先调用 runtime._sync_messages(host) 完成邮箱同步,将待投递消息写入 ctx.incoming_messages 和 ctx.incoming_envelopes。因此,当 WaitMessage.update() 被调用时,ctx.incoming_messages 已经是同步后的最新状态,无需额外的异步操作。

行为语义

def update(self) -> Status:
    incoming = getattr(self.ctx, "incoming_messages", None)
    if isinstance(incoming, list) and incoming:
        return Status.SUCCESS
    return Status.RUNNING

节点的逻辑极其简洁:有消息则 SUCCESS,无消息则 RUNNING。这种设计使其天然适配行为树的记忆序列(memory=True 的 Sequence):当 WaitMessage 返回 RUNNING 时,父序列会挂起等待,下一次 tick 时重新从该节点开始执行,直到消息到达后继续执行序列中的后续子节点。

ctx.incoming_messages 的数据来源

mailbox.sync_messages() 在同步消息时不仅设置 ctx.incoming_messages(消息对象列表),还设置 ctx.incoming_envelopes(信封对象列表,包含 seq、路由元数据等),同时将消息合并到状态管理器的 "messages" 键中(通过 StateMessageStore._append_many)。这使得下游节点不仅可以通过 WaitMessage 感知消息到达,还可以直接访问 ctx.incoming_messages 获取消息内容、发送者等详细信息。

WaitForever:永不终止的辅助节点

WaitForever 是 WaitMessage 的镜像节点——它始终返回 Status.RUNNING,永不自行终止。其主要用途是让一个 Agent 在完成初始化工作后保持存活状态,等待外部中断(如 preempt_agent 或 kill)。在需要 Agent 长期驻留等待动态消息的场景中,WaitForever 常被放在行为树末尾,确保 Agent 不会因树执行完毕而被标记为完成。

SwarmNode:基于技能的 Swarm Agent 节点

SwarmNode 继承自 SkillNode,是为 SwarmRole 定制的技能执行节点。它在 SkillNode 的基础上做了三项关键调整:

输出 Schema 抑制:_resolve_output_schema() 覆盖为返回 None(除非显式设置),因为 Swarm Agent 的最终输出通常是通过 send_message / send_direct_message 发送的纯文本,而非技能的 Pydantic 结构化输出。

输出键策略:_resolve_output_key() 优先使用显式设置的 output_key,否则使用节点名称——这避免了单技能场景下自动继承技能专用输出键的行为。

对话完成策略注入:_skill_react_system_prompt() 在基础 Prompt 后追加了"对话完成策略"指令,明确告知 LLM:一次完整的 turn 必须在调用 send_message 或 send_direct_message 发送可见回复后才算完成,避免 Agent 在内部推理或工具执行后提前终止。

此外,update() 在执行父类逻辑后,会将 swarm_success 标记写入状态管理器,供外部观察者(如 AgentTeam 或测试代码)判断 Agent 是否成功完成。

完整通信模式:Coordinator-Worker 协作

以下示例来自 basic_collaboration_demo.py,展示了三个核心节点协同工作的最小完整模式:

sequenceDiagram
    participant C as Coordinator Tree
    participant R as AgentRuntime
    participant W as Worker Tree

    C->>R: SpawnAgent("SpawnWorker", role="worker")
    R->>W: 创建 AgentHost + 启动
    R-->>C: agent_id → state["worker_id"]

    C->>C: PrepareTaskNode: 构造 task 消息
    C->>R: SendMessage(to_agent_key="worker_id", content_key="task")
    R->>W: 投递消息到 Worker inbox
    R->>W: mailbox sync → ctx.incoming_messages

    W->>W: WaitMessage → SUCCESS (有消息)
    W->>W: BuildWorkerReplyNode: 构造回复
    W->>R: SendMessage(to_agent_key="reply_to", content_key="reply")
    R->>C: 投递回复到 Coordinator inbox

    C->>C: WaitMessage → SUCCESS (有回复)
    C->>C: CaptureReplyNode: 提取 final_reply

Coordinator 行为树是一个带记忆的 Sequence:

Sequence(name="CoordinatorFlow", memory=True, children=[
    SpawnAgent("SpawnWorker", role="worker", output_key="worker_id",
               runtime=self.runtime, reuse_if_present=True),
    PrepareTaskNode(),
    SendMessage("SendTaskToWorker", to_agent_key="worker_id",
                content_key="task", runtime=self.runtime),
    WaitMessage("WaitWorkerReply"),
    CaptureReplyNode(),
])

Worker 行为树同样使用记忆序列:

Sequence(name="WorkerFlow", memory=True, children=[
    WaitMessage("WaitTask"),
    BuildWorkerReplyNode(),
    SendMessage("SendReply", to_agent_key="reply_to",
                content_key="reply", runtime=self.runtime),
])

reuse_if_present=True 的设置在 Coordinator 的行为树中至关重要:若行为树因某种原因被重新执行(例如记忆序列的回退),已有的 worker_id 不会被覆盖,从而避免创建冗余 Agent。

与 Swarm 运行时工具的互补关系

三个核心节点面向静态行为树编排——开发者在设计时明确 Agent 之间的通信拓扑。与之互补的是 SwarmToolProvider 提供的运行时工具集,它们允许 LLM 驱动的 Agent(通过 SwarmRole / SwarmNode)在运行时动态做出通信决策:

工具名 对应节点 差异
create / create_agent SpawnAgent LLM 动态决定何时创建、创建什么角色;支持 reuse_key 幂等
send_message / send_direct_message / send_topic_message SendMessage LLM 选择目标、内容;支持 dedupe_key 幂等
subscribe_topic 无直接节点对应 运行时动态加入主题,WaitMessage 随后可接收该主题消息
list_agents / self / pause_agent / resume_agent / kill_agent 无直接节点对应 提供 Agent 自省和生命周期管理能力

这种"静态节点 + 动态工具"的双层设计使 Jianmu 的多 Agent 系统既能支持确定性的协作流程(如 Coordinator-Worker),也能支持 LLM 自组织的液态拓扑(如 autonomous_swarm_demo.py 中的 Leader-Scout-Analyst 三级链式路由)。

阅读建议

至此,你已经掌握了构建多 Agent 通信的节点原语。要理解这些节点背后的完整运行时机制,建议继续阅读: