跳转至

Swarm runtime cn

AgentRuntime 是 Jianmu 多 Agent 协作的运行时引擎——它不像若干独立的单 Agent 循环那样各自为政,而是将消息路由、邮箱投递、事件总线和生命周期控制统一在同一个调度器之下。每当一个 Agent 通过 spawn() 诞生、通过 send_message() 发出消息、或在 auto 模式下持续 tick,背后都依赖本页所述的 Runtime 组件协同工作。

架构总览:从 AgentRuntime 看全局组件关系

AgentRuntime 是核心门面类(约 1160 行),其内部组合了六个关键子系统,各自承担明确职责。下图展示了它们的依赖方向与数据流:

flowchart TB
    subgraph Runtime["AgentRuntime (core.py)"]
        HOSTS["Host Registry<br/>_hosts: dict[str, AgentHost]"]
        EB["Event Bus<br/>InMemoryEventBus"]
        MS["Message Store<br/>InMemoryMessageStore"]
        RIDX["Routing Indexes<br/>_group_members / _topic_subscribers / _attachments"]
        TPS["Tool Providers<br/>[SwarmToolProvider, MCPToolProvider, ...user...]"]
        PR["Prompt Runtime<br/>_prompt_runtime"]
    end

    subgraph Scheduler["scheduler.py"]
        STEP["step_agent / step_all / run_until_idle"]
        AUTO["run_runner_task / ensure_runner_task"]
        RESTART["restart_policy + backoff"]
    end

    subgraph Mailbox["mailbox.py"]
        SYNC["sync_messages"]
        DRAIN["drain_mailbox"]
        PRIME["prime_mailbox_for_run"]
        GC["maybe_prune_message_store"]
    end

    subgraph Routing["routing.py"]
        SEND["send_message"]
        ENQ["enqueue_hot_envelope"]
        WAKE["wake_targets"]
        SUB["join_group / subscribe_topic"]
    end

    subgraph Host["AgentHost (host.py)"]
        PROFILE["profile + role"]
        RUNNER["ReactiveRunner"]
        STATE["StateManager"]
        INBOX["inbox: deque + pending_messages"]
    end

    Runtime --> Scheduler
    Runtime --> Mailbox
    Runtime --> Routing
    HOSTS --> Host
    Scheduler --> Host
    Mailbox --> MS
    Mailbox --> Host
    Routing --> MS
    Routing --> HOSTS
    EB -->|"emit(AgentEvent)"| HOSTS

AgentRuntime.__init__ 接收约 20 个可配置参数,涵盖角色注册、消息存储、工具链、重启策略、快照检查点等全部维度。在构造阶段它会按顺序排列 ToolProvider 链:内置 SwarmToolProvider → 可选的 MCPToolProvider → 用户传入的额外 Provider,确保下游可覆盖上游同名工具。

AgentHost:每个 Agent 的运行时载体

AgentHost 是一个轻量级数据类(dataclass),它将单个 Agent 的所有运行时状态封装在一起,是调度器、邮箱和路由系统操作的原子单位。

字段 类型 用途
profile AgentProfile Agent 标识(id、role、task、parent_id)
role AgentRole 用于构建行为树的角色工厂
runner ReactiveRunner 事件驱动的异步 tick 执行器
state StateManager 类型化状态管理
mailbox_cursor_seq int 已消费的消息序列号游标
groups set[str] 当前订阅的广播组
topics set[str] 当前订阅的主题
inbox deque[MessageEnvelope] 热路径直达消息队列
pending_messages list[MessageEnvelope] 待处理的消息列表
paused / auto_mode bool 生命周期标志位
task / sync_task asyncio.Task 后台 runner 任务和邮箱同步任务
mailbox_dirty / mailbox_force_signal / mailbox_version — 邮箱变更检测机制
side_effects dict 工具副作用去重(如 idempotent send/create)

mailbox_version 是一个单调递增计数器,每次 mark_dirty() 调用都会使其递增。调度器在执行 runner 任务前后通过比较版本号来判断是否出现了新消息——这就是"新消息驱动继续运行"的检测机制。

消息路由:三种投递模式与热路径优化

消息从发送到被目标 Agent 消费,经历三个层次的处理管道:

sequenceDiagram
    participant Sender as 发送方 Agent
    participant Routing as routing.send_message
    participant Store as InMemoryMessageStore
    participant Inbox as 目标 Host.inbox
    participant Mailbox as mailbox.sync_messages
    participant State as StateMessageStore

    Sender->>Routing: send_message(content, to_agent_id/group_id/topic)
    Routing->>Store: append(MessageEnvelope) → seq
    Routing->>Inbox: enqueue_hot_envelope(env)
    Routing->>Sender: wake_targets()
    Note over Sender: emit("message_sent")

    Mailbox->>Inbox: 检查 inbox deque
    Mailbox->>Store: read_since(cursor, agent_id, groups, topics)
    Mailbox->>State: _append_many(incoming_messages)
    Note over Mailbox: 可选 context_builder.transform 合并历史
    Mailbox->>Store: ack(seq) 确认消费

三种路由目标

send_message 要求恰好指定一个目标(to_agent_id、group_id 或 topic),否则抛出 ValueError。三种模式的对比如下:

模式 目标参数 投递语义 适用场景
直连 (Direct) to_agent_id 精确送达指定 Agent 点对点对话、任务委派
组播 (Group) group_id 送达组内所有成员 广播通知、团队同步
主题 (Topic) topic 送达所有订阅者 发布/订阅模式、事件驱动

热路径优化:enqueue_hot_envelope

消息存储到 InMemoryMessageStore 后,并不等到 Agent 下一次 mailbox sync 时才被发现——enqueue_hot_envelope 会直接将 MessageEnvelope 推入所有匹配目标 Agent 的 host.inbox deque。这意味着在同一个事件循环迭代中,目标 Agent 就能感知到新消息,大幅降低延迟。

运行时本地生产者必须使用 AgentRuntime.enqueue_internal_message(),而不是直接修改 host.inbox。这一统一交接会推进 mailbox version、记录可选的 mailbox_source,并可选地唤醒 Host,因此调度器不会遗漏本地注入的输入。

唤醒机制:wake_targets

消息发送后,wake_targets 负责通知接收方"有活干了"。直连模式下直接 runtime.wake(agent_id);组播/主题模式下遍历索引中的成员并逐一唤醒。它还支持 suppress_group_wake 元数据标记和显式 wake_agent_ids 列表,允许发送方精确控制唤醒范围。

组与主题管理

join_group / leave_group 和 subscribe_topic / unsubscribe_topic 操作 Agent 在 Runtime 级索引中的成员关系,同时发出对应的 group_joined、topic_subscribed 等事件。当 Agent 被 kill 时,remove_host_routing 会清理其在所有索引中的残留。

邮箱系统:消息同步与游标消费

邮箱系统在 mailbox.py 中实现了三个核心操作,它们共同构成了 Agent 消息消费的完整生命周期:

1. sync_messages:将外部消息载入 Agent 状态

这是消息消费的主入口。它按优先级从两个来源获取消息: - 优先检查 host.inbox(热路径推送的消息),避免额外的存储查询 - 回退到 message_store.read_since(cursor) 从持久化存储读取

读取到的消息通过 StateMessageStore._append_many 写入 Agent 的状态管理器中的 messages 字段。如果 Runtime 配置了 context_builder,还会在合并历史消息时调用 context_builder.transform() 进行处理(例如 Token 预算控制、消息去重)。

每个成功同步的消息都会发出 message_received 事件,并推进 host.mailbox_cursor_seq 游标。

2. prime_mailbox_for_run:Runner 启动前同步

在每次 runner 执行前调用。它处理当前可能正在进行的 sync_task(等待其完成或清理),然后通过 sync_lock 获取锁进行同步。同步完成后,如果仍然 mailbox_dirty,会重新调度 schedule_mailbox_sync。

drain 到 prime 的交接

drain_mailbox 和 prime_mailbox_for_run 的执行顺序并不固定。若 drain 在新建 runner 到达 prime_mailbox_for_run 前已同步一批消息,会将该 incoming context 标记为已准备;即使 dirty 标志已清除,prime_mailbox_for_run 仍会为该次 execution 领取这批消息。未被 drain 明确标记为已准备的 context 则属于上一轮的陈旧上下文,会在新 execution 前清理。这样既避免吞掉新消息,也避免复用前一轮 execution 的 incoming batch。

3. drain_mailbox:有界迭代邮箱排空

在 auto 模式下,后台 sync_task 循环调用此函数。它在 mailbox_drain_max_iterations(默认 100)次迭代中反复检查 mailbox_dirty 标志并同步消息。每次迭代后都检查是否需要继续;达到上限时发出 mailbox_drain_capped 警告并调度后续继续。

死信处理

on_sync_failure 实现了消息投递的失败重试与死信队列机制。对于每次同步失败的 pending_messages,调用 message_store.fail() 累加失败计数;超过 delivery_max_failures(默认 3)阈值的消息通过 mark_dead() 移入死信队列,并发出 message_dead_letter 事件。

消息存储 GC

maybe_prune_message_store 在所有 Host 游标均已推进时触发。它计算所有 Host 游标的最小值作为安全删除边界,并确保两次 GC 之间的序列号跨度超过 message_store_gc_min_advance(默认 1000),以防止频繁回收。

InMemoryEventBus:线程安全的事件总线

InMemoryEventBus 是 Swarm 运行时内部事件的发布/订阅中枢。它基于环形缓冲区 + 订阅者列表设计,所有操作受 threading.Lock 保护。

classDiagram
    class InMemoryEventBus {
        -_subscribers: List[Callable]
        -_buffer: List[Tuple[float, AgentEvent]]
        -_max_buffer: int
        -_lock: Lock
        +emit(event) None
        +subscribe(callback) Callable
        +get_since(since) List[AgentEvent]
        +get_all() List[AgentEvent]
    }
    class AgentEvent {
        +event_type: str
        +agent_id: str | None
        +payload: dict
        +created_at: float | None
        +event_id: str | None
    }
    InMemoryEventBus --> AgentEvent

emit() 方法自动补全 event_id(UUID4)和 created_at(Unix 时间戳),然后将事件追加到缓冲区并通知所有订阅者。单个订阅者的异常不会影响其他订阅者。订阅通过返回 unsubscribe 回调来管理:调用该回调即从订阅列表中移除。

AgentRuntime.emit() 在调用 event_bus.emit() 之前会做一次规范化——将 session_id 和 agent_id 注入 payload——然后桥接选定的生命周期事件到 Runtime 级 RuntimeEventBus(agent.created、agent.started 等)。

调度器:手动模式与自动模式的双轨制

Swarm Runtime 提供两种 Agent 运行模式,分别对应不同场景:

维度 手动模式 (Manual) 自动模式 (Auto)
触发方式 显式调用 step_agent() 等 API spawn() 时 auto_start=True 或 start_agent()
执行循环 调用方控制每步 tick 后台 asyncio.Task 持续 runner.run()
并发控制 step_all_concurrent(max_concurrency=N) 各自独立 task 并发
适用场景 调试、交互式步进、外部编排 全自动多 Agent 协作

手动模式 API

  • step_agent(agent_id):推进单个 Agent 一个 tick,返回 Status。内部通过 sync_lock 完成消息同步后调用 runner.step()。
  • step_agent_with_actions(agent_id):同上,但额外返回 actions 字典。
  • step_all():按主机注册顺序依次推进所有 Agent。
  • step_all_concurrent(max_concurrency=N):通过 asyncio.Semaphore 并发推进所有 Agent。
  • run_until_idle(max_steps=200):持续推进直到所有 Agent 处于空闲状态(无新消息、无 RUNNING 节点、无活跃 async_task)。

自动模式核心循环:run_runner_task

这是 auto 模式下每个 Host 的后台主循环,体现了 Runtime 的完整重启与唤醒语义:

flowchart TD
    START["runner task 启动"] --> RESET["reset_runner_for_restart()"]
    RESET --> PRIME["prime_mailbox_for_run()"]
    PRIME --> RUN["runner.run()"]
    RUN --> CHECK_RETRY{"收到消息但未发送?<br/>same_run_no_send_retries"}
    CHECK_RETRY -->|是| RESET
    CHECK_RETRY -->|否| CHECK_CONT{"新消息到达?<br/>should_continue_after_run"}
    CHECK_CONT -->|是| RESET
    CHECK_CONT -->|否| EMIT_DONE["emit agent_run_completed"]

    RUN -->|异常| ERR_HANDLER{"restart_policy?"}
    ERR_HANDLER -->|never| FAILED["emit agent_failed"]
    ERR_HANDLER -->|on_failure/always| CHECK_RESTART{"restart_count < max_restarts?"}
    CHECK_RESTART -->|是| BACKOFF["backoff 延迟"]
    BACKOFF --> RESET
    CHECK_RESTART -->|否| FAILED

should_continue_after_run 通过比较 mailbox_version 是否变化、检查 mailbox_dirty、以及调用 message_store.has_unread() 来判断运行期间是否有新消息到达,实现了消息驱动的连续运行。

should_retry_same_run_without_send 处理一种边缘情况:Agent 收到消息并完成了一轮运行,但在该轮中没有发送任何消息。此时如果 same_run_no_send_retries 允许,运行体会自动重试——确保"收到消息后一定回复"的语义。

空闲检测

is_idle 要求所有 Host 同时满足:mailbox_dirty 为 False、无活跃 sync_task、根节点状态非 RUNNING、且 has_active_work 返回 False。has_active_work 不仅检查消息队列,还遍历行为树中的所有节点:检查 has_pending_work 回调、活跃的 async_task、RUNNING 状态的 AsyncBehaviour。

Agent 生命周期:从创建到销毁的完整状态机

stateDiagram-v2
    [*] --> Created: spawn(role, task)
    Created --> Running: auto_start 或 start_agent()
    Created --> Stepping: step_agent() (manual)

    Running --> Paused: pause()
    Paused --> Running: resume()

    Running --> Preempted: preempt_agent()
    Preempted --> Running: 自动 restart

    Running --> Error: runner 异常
    Error --> Running: restart (backoff)
    Error --> Failed: restart 超限或 policy=never

    Running --> Stopped: kill() / kill_async()
    Stepping --> Stopped: kill() / kill_async()
    Paused --> Stopped: kill() / kill_async()
    Failed --> Stopped: kill() / kill_async()

    Stopped --> [*]

spawn:从 Role 到 Host 的装配

spawn() 依次完成以下步骤: 1. 解析 role(字符串查找注册表或直接使用对象) 2. 生成 UUID agent_id,构建 AgentProfile 3. 处理 constraints(反序列化 → 预算规格提取 → 元数据注入) 4. 调用 role.build_tree(profile) 构建行为树根节点 5. 通过 state_factory 创建 StateManager 并初始化 6. 组装 RunContext(注入 model_client、runtime、sandbox、prompt_runtime) 7. 通过 runner_factory 创建 ReactiveRunner 8. 创建 AgentHost,处理 attached 模式父子关系 9. 注册到 _hosts 字典,发出 agent_created 事件 10. 若 auto_start=True,立即调用 start_agent()

kill:级联终止

kill() 通过 _detach_host 从 Runtime 中移除 Host,同时返回 attached 子 Agent 列表。然后对子 Agent 递归调用 kill(),实现级联终止。最后将分离的 sync_task 和 task 加入后台清理队列。

kill_async() 与之功能相同,但使用 await 等待任务取消完成,适用于需要确保清理完毕的异步上下文。

pause / resume / wake

  • pause:设置 host.paused = True,发出 agent_paused 事件。已暂停的 Host 不再被调度器处理。
  • resume:设置 host.paused = False,发出 agent_resumed,然后调用 wake() 触发工作。
  • wake:标记 mailbox dirty + 强制信号 + ensure_runner_task,确保 Host 的 runner 任务存在并开始处理。

preempt_agent:中断并重启

preempt_agent 用于打断正在运行的 Agent 并强制重新开始。它取消当前 sync_task 和 task,通过 reset_runner_for_restart 重置 runner 状态(停止所有节点、取消异步任务),然后启动替代 runner task。替代 task 负责唯一一次 prime_mailbox_for_run,避免 mailbox drain 已准备的消息被重复领取或清除。

工具系统:Swarm 能力以 Tool 形式暴露

SwarmToolProvider 将 Runtime 的核心管理操作包装为 LLM 可调用的 Tool,通过 get_all_tools() 统一暴露:

Tool 名称 功能 关键参数
create_agent / create 创建子 Agent role, guidance, reuse_key
send_message 发送路由消息 to, content, group_id, topic, dedupe_key
send_direct_message 直连消息 to_agent_id, content, dedupe_key
send_topic_message 主题广播 topic, content, dedupe_key
list_agents 列出活动 Agent —
self 返回当前 Agent 档案 —
pause_agent 暂停 Agent agent_id
resume_agent 恢复 Agent agent_id
kill_agent 终止 Agent agent_id

这些工具内置幂等性支持:create_agent 使用 reuse_key 和 host.side_effects 字典确保相同 reuse_key 的重复调用返回已创建的 agent_id;send_message 使用 dedupe_key 和 effect_tags=("message_send",) 防止重复发送。

工具链采用 "后注册者覆盖同名工具" 的合并策略(ToolSet.from_tools(merged).with_tools(named_tools)),这使得用户可以注册自定义 Provider 来覆盖内置工具。

节点集成:SpawnAgent、SendMessage、WaitMessage

三个行为树节点将 Runtime 能力嵌入到 Agent 的执行树中:

  • SpawnAgent:异步节点。通过 _resolve_runtime 查找 Runtime 引用(支持属性注入、状态查找和全局回退),调用 runtime.spawn() 并将新的 agent_id 写入 output_key。支持 reuse_if_present 模式避免重复创建。

  • SendMessage:异步节点。从状态中解析 content、to_agent_id/group_id/topic 和 sender_id,调用 runtime.send_message()。当 Runtime 不可用时降级为通过 runtime.emit() 发送事件。

  • WaitMessage:同步节点。仅检查 ctx.incoming_messages 是否非空——这是 mailbox.sync_messages 注入的上下文属性。有空消息时返回 SUCCESS 让行为树推进,否则返回 RUNNING 保持阻塞。

SwarmNode 继承自 SkillNode,增加了 swarm 特有的对话完成策略提示("必须调用 send_message 才算完成一轮"),以及通过 swarm_success 状态标志位记录执行结果。

快照与恢复:运行时持久化

save_snapshot 将整个 Runtime 状态序列化:

{
    "version": 1,
    "session_id": "...",
    "hosts": [
        {
            "profile": {...},        # AgentProfile 字典
            "state": {...},          # StateManager checkpoint
            "mailbox_cursor_seq": N,
            "groups": [...],
            "topics": [...],
            "attached_to": "...",
            "paused": bool,
            "auto_mode": bool,
            ...
        }
    ],
    "message_store": {...}          # InMemoryMessageStore dump_state()
}

restore_snapshot 执行反向流程:先 kill 掉当前所有 Host → 恢复消息存储 → 逐个重建 Host(包括重新 build_tree、恢复状态、重建 RunContext 和 ReactiveRunner)→ 重建路由索引 → 重启 auto 模式任务。

序列化通过 CheckpointerProtocol 的 save_checkpoint/get_checkpoint 接口完成,当前支持 FileCheckpointer(文件系统后端)。

可观测性:事件流与桥接

Runtime 提供两层事件订阅:

  1. subscribe_events(callback):订阅 InMemoryEventBus 上所有 AgentEvent。每次事件在 publish 前经过 emit() 规范化(注入 session_id 和 agent_id)。

  2. subscribe_observability(callback, runtime_events=True, trace_events=False):统一的观测入口。runtime_events=True 转发规范化的 Runtime 事件;trace_events=True 额外订阅 telemetry trace 事件并以统一格式转发。

启用 trace 转发时,订阅会跟随进程级默认 TelemetryHub 的替换;telemetry bootstrap 或其他 Hub 替换后,调用方无需重新订阅。

Runtime 的完整事件类型清单涵盖:agent_created、agent_started、agent_restarted、agent_paused、agent_resumed、agent_wake、agent_preempted、agent_stopped、agent_killed、agent_failed、agent_error、agent_step_completed、agent_run_completed、agent_mailbox_continue、agent_no_send_after_incoming、agent_restart_scheduled、message_sent、message_received、message_dead_letter、mailbox_drained、mailbox_drain_capped、mailbox_sync_error、message_store_compacted、group_joined、group_left、topic_subscribed、topic_unsubscribed、runtime_restored。

其中选定的生命周期事件会桥接到 RuntimeEventBus(如 agent.created),使得上层 Agent 门面和遥测系统能够统一感知 Swarm 生命周期。

高层组装:SwarmRole 与 AgentTeam

SwarmRole(role.py)是面向"液态拓扑"协作场景的预设角色实现。它在 build_tree() 中构造一个 SwarmNode,内部采用 ReAct 循环 + 技能系统,将 system_prompt 和 task 拼接为完整 prompt,并通过 _resolve_runtime_tools() 获取 Runtime 工具。

AgentTeam(team.py)是 AgentRuntime 的外观包装器,负责: - 从项目配置 jianmu.yaml 中读取 checkpoint 默认值 - 统一管理 defaults(model_client、context_builder、tool_providers) - 提供 async with 上下文管理器接口(__aenter__ → initialize_async,__aexit__ → close_async)

两者在架构中的关系是:AgentTeam → AgentRuntime → AgentHost → ReactiveRunner(行为树)→ SwarmNode → SwarmRole(配置)。

阅读建议

本文档覆盖了 Swarm 运行时的内部引擎机制。接下来建议阅读: - 角色与团队:AgentRole、SwarmRole 与 AgentTeam 的协作编排 —— 深入了解角色工厂、@role 装饰器和团队级配置 - Swarm 节点:SpawnAgent、SendMessage、WaitMessage 的多 Agent 通信原语 —— 从节点视角理解 Swarm 原语在行为树中的用法 - ReactiveRunner:事件驱动的异步 tick 调度与挂起恢复机制 —— AgentHost 底层的执行器详解