跳转至

Runtime cn

ReactiveRunner 是 Jianmu 运行时引擎的核心调度器。它封装了 py_trees 行为树的标准执行模型,在其之上构建了一套事件驱动的异步 tick 循环——节点完成异步任务后通过唤醒信号通知调度器,而不是让调度器盲目轮询;当工作流需要外部交互(审批、用户输入、外部执行)时,调度器可以将树挂起(suspend)并归还控制权给宿主应用,待外部条件满足后恢复(resume)继续执行。三者共同构成 Jianmu "运行时永不停机但可优雅挂起"的执行语义。

架构全景

ReactiveRunner 位于行为树执行内核与上层 Agent 外观之间,向下通过 BehaviourTree.tick() 驱动节点图,向上通过 Agent / AgentTeam 外观暴露 run()、step()、run_until_suspend() 等高阶入口。

flowchart TB
    subgraph Host["宿主应用层"]
        Agent["Agent / AgentTeam 外观"]
    end

    subgraph Runner["ReactiveRunner"]
        direction TB
        Entry["run() / step() / run_until_suspend()"]
        Loop["_event_loop() 事件循环"]
        Signal["tick_signal (asyncio.Event)"]
        Wake["_wait_for_tick() / _signal_tick()"]
        Tick["tick_once() → tree.tick()"]
    end

    subgraph Core["运行时核心"]
        SM["StateManager"]
        IC["InteractionController"]
        BT["py_trees BehaviourTree"]
    end

    subgraph Nodes["节点层"]
        AsyncNode["AsyncBehaviour"]
        SyncNode["Behaviour / Node"]
    end

    Agent --> Entry
    Entry --> Loop
    Loop --> Wake --> Signal --> Loop
    Loop --> Tick --> BT
    BT --> AsyncNode
    BT --> SyncNode
    AsyncNode -->|"wake_up() 回调"| Signal
    SM -->|"subscribe / notify"| Signal
    IC -->|"suspension 检测"| Loop

ReactiveRunner 在构造阶段完成三项关键初始化:(1) 将 RunContext(模型客户端、审批管理器、沙箱句柄等)和 StateManager 通过依赖注入传递给树中每一个 JianmuNodeMixin 节点;(2) 调用 tree.setup() 完成 py_trees 标准初始化;(3) 注册 StateManager 的变更订阅回调,使得任何状态写入都能唤醒事件循环。

三层执行 API:从原子 tick 到全自动循环

ReactiveRunner 提供三个粒度的执行入口,满足从调试到生产的多种场景。

API 语义 异步 模式锁定 适用场景
tick_once() 执行一次同步 tree.tick(),返回根节点状态 否 无 最低层 API,单元测试
step() 一次异步 tick + step-scoped 状态重置 是 限制 step/run 互斥 TUI/Studio 的交互式步进
run() 事件驱动循环直至 SUCCESS/FAILURE 是 限制 step/run 互斥 生产级端到端 Agent 运行

step() 与 run() 通过内部 _mode 标志("idle" | "step" | "run")实现互斥——不能在 run() 进行中调用 step(),反之亦然,避免了并发 tick 导致的状态撕裂。

tick_once:原子 tick

tick_once() 是推理内核的最细粒度操作——它直接调用 self.tree.tick(),令 py_trees 从根节点向下遍历一次整棵树,完成后返回根节点状态。它不关心事件循环,不触发任何 suspend/resume 逻辑,也不参与状态管理。

runner.tick_once()  # 一次同步 tick,返回 Status.SUCCESS / FAILURE / RUNNING

step() 方法正是基于 tick_once() 构建的:先调用 state_manager.reset_step_fields() 清空 Ephemeral("step") 标记的临时字段,再合并观察数据,然后执行一次 tick,最后返回 step-scoped 字段的快照。

step:交互式步进

step() 专为 TUI Chat 和 Tree Studio 等需要逐 tick 外部控制的交互式应用设计。它在一次 tick 前后完成 step-scoped 状态管理,并将本轮产生的动作字段(如 final_answer)返回给调用方,使得 UI 可以在每个 tick 之间插入渲染逻辑。

# TUI 应用中的典型用法
result = await runner.step(obs={"user_message": "hello"})
# result = {"final_answer": "你好!有什么可以帮你?"}

调用方通过 yield_to_async=True 可以在 tick 后额外让出事件循环,确保 UI 框架有机会处理挂起的渲染事件。

run 与 _event_loop:全自动事件驱动循环

run() 是生产环境的标准入口。它内部委托给 _event_loop(),后者实现一个事件驱动而非忙等的 tick 循环:

sequenceDiagram
    participant Host as 宿主
    participant Loop as _event_loop()
    participant Signal as asyncio.Event
    participant Tree as BehaviourTree
    participant Node as AsyncBehaviour

    Host->>Loop: run(input_data)
    Loop->>Loop: 初始化、注入依赖、kick 首个 tick
    Loop->>Signal: _signal_tick() (首次)
    loop 事件驱动循环
        Loop->>Signal: _wait_for_tick() 阻塞等待
        Note over Loop,Signal: 不会忙等,await Event.wait()
        Loop->>Tree: tree.tick()
        Tree->>Node: tick() → update_async()
        Node-->>Tree: 返回 RUNNING (任务未完成)
        Tree-->>Loop: Status.RUNNING
        Loop->>Loop: 速率限制 (max_fps)
        Note over Loop,Node: 异步任务完成后...
        Node->>Signal: wake_up() → _signal_tick()
        Signal-->>Loop: Event.set() 唤醒
        Loop->>Tree: tree.tick() (下一轮)
        Tree->>Node: 收集已完成任务结果
        Tree-->>Loop: Status.SUCCESS / FAILURE
    end
    Loop-->>Host: 最终 Status

核心机制:循环在每次 tick 后通过 _wait_for_tick() 进入阻塞等待——不是忙等,而是 await asyncio.Event.wait()。当以下任一事件发生时,tick_signal 被置位,循环被唤醒:

  1. 节点异步任务完成:AsyncBehaviour.initialise() 在启动 update_async() 协程时,通过 add_done_callback(lambda _: wake_up()) 注册完成回调,任务结束时自动通知 runner。
  2. StateManager 状态变更:runner 在构造时订阅了 StateManager 的变更通知,任何节点的 write_state() 调用都会间接触发 _on_wake_signal()。
  3. 审批决议到达:当 SuspensionMode.WAIT 模式下审批被外部决议时,live approval 订阅回调会调用 _signal_tick()。

唤醒信号合并:防止事件风暴

在高频场景(如工具密集调用、多节点同时完成)中,可能短时间内产生大量唤醒信号。ReactiveRunner 通过 _pending_tick_count 计数器实现唤醒合并:

tick_signal 被 set 时:
  if pending_tick_count < max_pending_wakeups:
      pending_tick_count += 1

_wait_for_tick() 被唤醒时:
  if pending_tick_count > 0:
      pending_tick_count -= 1
      if pending_tick_count == 0:
          tick_signal.clear()
      return  # 消费一个唤醒信号
  tick_signal.clear()
  await tick_signal.wait()  # 无信号则继续等待

默认 max_pending_wakeups = 2,意味着最多缓冲 2 个待处理唤醒信号——1 个正在飞行中 + 1 个额外边缘触发。超过上限的信号被静默丢弃,防止无界积压。这在 _wait_for_tick() 的设计注释中被描述为 "advisory signals rather than transactional events"——唤醒信号是建议性的,不是事务性的必须逐条消费的消息。

线程安全:跨线程唤醒

由于 StateManager 可能被非 asyncio 线程写入(例如来自外部 HTTP 回调或信号处理器),_signal_tick() 必须处理跨线程场景:

def _signal_tick(self):
    loop = self._loop
    if loop is None: return        # 事件循环已关闭,丢弃
    with self._tick_lock:
        if self._pending_tick_count < self._max_pending_ticks:
            self._pending_tick_count += 1
    if threading.get_ident() == self._loop_thread:
        self.tick_signal.set()      # 同一线程,直接设置
        return
    try:
        loop.call_soon_threadsafe(self.tick_signal.set)  # 跨线程安全
    except RuntimeError:
        ...  # 事件循环已关闭,丢弃

_event_loop() 在进入时捕获当前事件循环和线程 ID,退出时清理。如果 call_soon_threadsafe 因循环已关闭而抛出 RuntimeError,信号被静默丢弃——这在实际运行中通常意味着 runner 已正常退出。

挂起与恢复:交互式工作流的中断-继续机制

当行为树中的节点需要外部交互(审批确认、用户输入、外部系统回调)时,节点通过 InteractionController 创建 SuspensionRecord 写入 runtime metadata。ReactiveRunner 在每个 tick 后检查是否存在活跃挂起,并根据 SuspensionMode 决定行为。

挂起分类与运行时语义

SuspensionCategory SuspensionReason 典型触发场景 恢复方式
REQUIRE_APPROVAL APPROVAL_PENDING Guard 体系拦截工具调用 审批管理器决议 + execution.resumed 事件
REQUIRE_USER_INPUT AWAITING_USER_INPUT 节点请求用户补充信息 宿主注入 resume payload
REQUIRE_EXTERNAL_EXECUTION EXTERNAL_RESULT_PENDING 委托外部系统执行 宿主回调写入执行结果

三种挂起类别共享同一套 InteractionController → StateManager → runtime_metadata 持久化路径,区别仅在于恢复时的 payload 规范化方式不同。

SuspensionMode:WAIT vs YIELD

flowchart TD
    Tick["tree.tick() 完成"] --> Check{"存在活跃 Suspension?"}
    Check -->|否| Continue["继续循环"]
    Check -->|是| Mode{"SuspensionMode?"}
    Mode -->|WAIT| Subscribe["保持循环运行\n订阅 live approval 回调"]
    Subscribe --> Wait["继续等待 tick_signal\n(审批决议会触发唤醒)"]
    Mode -->|YIELD| Save["强制保存 checkpoint"]
    Save --> Return["return RunResult(outcome='suspended')\n控制权归还宿主"]
  • WAIT 模式:runner 保持事件循环运行,等待 ApprovalManager 的 live 决议回调。适用于审批场景——GUI 弹出确认对话框的同时,后台循环保持存活,一旦用户点击批准,回调自动触发唤醒。
  • YIELD 模式:runner 立即退出事件循环,保存 checkpoint,返回 RunResult(outcome="suspended", suspension=...) 给宿主。宿主可以在任意时刻通过 resume_and_continue() 恢复执行。适用于 Web 服务或无头批处理场景。

恢复流程:run_until_suspend 与 resume_and_continue

run_until_suspend() 是 run() 的扩展版本,它在事件循环中增加了挂起检测和 RunResult 返回路径。方法签名支持多种恢复策略:

参数 类型 作用
restore RestorePolicy NEVER:全新运行;IF_EXISTS:有 checkpoint 则恢复;REQUIRED:必须从 checkpoint 恢复
resume_data dict 宿主提供的恢复 payload(用户输入值 / 外部执行结果)
resume_request_id str 必须与活跃 suspension 的 request_id 精确匹配
suspension_mode SuspensionMode 本次运行的挂起处理策略

恢复流程的核心验证链:

resume_request_id → build_resume_data() → raise_for_resume_result()
                         ↓
              SuspensionRecord 匹配检查
                         ↓
           APPROVAL 类 → resolve_approval_resume() → approval_manager
           其他类别     → write_resume_payload()    → StateManager
                         ↓
              deactivate_suspension()
                         ↓
              发出 execution.resumed 事件

resume_and_continue() 是对 run_until_suspend(restore=RestorePolicy.REQUIRED, ...) 的便利封装,专用于"从已保存 checkpoint 恢复挂起的工作流并继续执行"的标准路径。

Live Approval:WAIT 模式下的实时决议桥接

当 SuspensionMode.WAIT 下出现 REQUIRE_APPROVAL 挂起时,ReactiveRunner 通过 _sync_live_approval_subscription() 在 ApprovalManager 上注册一个一次性决议回调。该回调在审批被外部决议时触发,执行以下桥接:

  1. 通过 _bridge_live_approval_resolution() 验证决议与当前活跃 suspension 匹配
  2. 将决议写入 runtime metadata 的 resolved_approvals inbox
  3. 调用 deactivate_suspension() 标记挂起已解除
  4. 发出 execution.resumed 运行时事件
  5. 调用 _signal_tick() 唤醒事件循环

整个过程在 loop.call_soon_threadsafe() 的保护下跨线程安全执行,确保即使审批决议来自 GUI 线程也能正确调度。

Checkpoint 持久化:状态与树拓扑的双重保存

ReactiveRunner 的 checkpoint 机制保存两个维度的数据:

  1. StateManager 全量状态:通过 dump_checkpoint() 导出 Pydantic 模型的所有字段和 runtime_metadata
  2. 行为树拓扑状态:每个节点的 Status 按路径键(path:0.1.2)序列化

路径键通过 _build_path_maps() 在运行时构建,它从根节点开始 DFS 遍历,为每个节点分配形如 "path:0.1.2" 的稳定标识——根为 "0",其第二个子节点为 "0.1",依此类推。这种路径方案比节点 ID 更稳定,因为即使节点对象在多次运行中重新创建,只要树结构不变,路径键就不变。

恢复时的 Composite 指针修复

py_trees 的 Composite 节点(Sequence、Selector)在被标记为 RUNNING 时,依赖 current_child 指针记住当前正在执行的子节点。从 checkpoint 恢复时,仅仅回放 Status.RUNNING 到 composite 是不够的——还必须修复 current_child 指向正确的子节点。

_repair_running_composite_pointers() 按 composite 类型差异化处理:

Composite 类型 current_child 修复规则
Sequence 指向第一个 status != SUCCESS 的子节点
Selector 指向第一个 status != FAILURE 的子节点
其他 指向第一个 INVALID 或 RUNNING 的子节点

如果找不到合适的子节点,则将 composite 整体标记为 INVALID,让 py_trees 在下一个 tick 重新进入。

检查点保存节奏

flowchart LR
    Tick["每个 tick 后"] --> Periodic{"total_tick_count % interval == 0?"}
    Periodic -->|是| Save["_maybe_save_checkpoint()"]
    Periodic -->|否| Suspend{"挂起且 mode=YIELD?"}
    Suspend -->|是| Force["_save_checkpoint_now() 强制保存"]
    Suspend -->|否| Skip["跳过"]

周期性保存由 checkpoint_interval 控制(默认每 tick 一次),而强制保存仅在 run_until_suspend() 的 YIELD 路径触发——确保宿主在恢复时有精确的挂起点快照。

热循环检测与速率限制

ReactiveRunner 实现了两级保护防止异常节点导致 CPU 空转:

速率限制:max_fps(默认 60)决定最小 tick 间隔 1.0 / max_fps ≈ 16.7ms。每个 tick 后强制 await asyncio.sleep(min_tick_interval - tick_elapsed),确保即使 tick 极快完成也不会超过频率上限。

热循环检测:当 1 秒内的 tick 次数超过 max_fps * hot_loop_warn_factor(默认 60 × 1.5 = 90)时,发出 WARNING 日志。这通常意味着节点在极短时间内反复从 RUNNING 回到 RUNNING(例如无 sleep 的忙等循环),提示可能存在逻辑错误。检测在每 1 秒窗口后重置计数器。

依赖注入生命周期

ReactiveRunner 在每个 run() / run_until_suspend() 的执行周期内管理依赖注入的完整生命周期:

进入 _event_loop():
  1. InjectPayload(wake_up=_on_wake_signal) → 注入所有节点
  2. _bind_callbacks() → 订阅 StateManager 变更
  3. 事件循环运行...

退出 _event_loop() (finally 块):
  1. auto_driving = False
  2. 清空 pending_tick_count
  3. 取消 StateManager 订阅
  4. InjectPayload(wake_up=None) → 清除所有节点的唤醒回调
  5. tree.interrupt() → 中断树执行

第 4 步的 wake_up=None 重注入至关重要——它确保 runner 退出后,节点持有的回调引用被清除,防止悬空引用导致的内存泄漏。每个 run 周期开始时(步骤 1)重新注入有效回调。

与上层模块的协作关系

ReactiveRunner 不直接面向最终用户——它通过 Agent 和 AgentTeam 外观暴露能力。这些外观负责:

  • 从 jianmu.yaml 读取 RuntimeConfig(max_fps、max_pending_wakeups、hot_loop_warn_factor)
  • 构建默认 FileCheckpointer(当 runtime.checkpoint.enabled = true)
  • 将 thread_id、checkpoint_interval 等配置统一传递给 runner

Agent 还额外提供 resume_interaction() 方法,将 ReactiveRunner.resume_and_continue() 的 ResumeError 异常转换为结构化的 ResumeInteractionResult,使 UI 层可以用返回值而非 try/except 处理恢复失败。

总结

ReactiveRunner 在 py_trees 的同步 tick 模型之上构建了三层关键扩展:(1) 基于 asyncio.Event 的事件驱动唤醒机制,消除了忙等轮询;(2) 跨线程安全的信号合并,防止高频场景下的事件风暴;(3) 通过 InteractionController 和 checkpoint 持久化实现的挂起-恢复协议,使 Jianmu 工作流可以在需要外部交互时优雅中断并在条件满足时精确续接。这三层扩展共同构成了 Jianmu 区别于传统行为树引擎的核心竞争力——一个既可以全自动运行、又可以在任意时刻挂起等待人类决策的混合执行模型。


建议阅读路径:本文聚焦于 ReactiveRunner 的调度与挂起恢复机制。要理解运行时的完整组装流程,请阅读 Agent 与 AgentTeam 外观模式;要了解挂起恢复在 Guard 审批中的具体应用,请阅读 人机交互挂起恢复;要深入理解状态管理的 Ephemeral 字段与 Reducer 合并,请阅读 类型化状态管理。