跳转至

Custom nodes cn

Jianmu 提供了两条创建自定义节点的路径:类继承路径——从 Node(同步)或 AsyncNode(异步)基类派生,覆写核心方法 update() 或 update_async();装饰器路径——用 @node 将一个纯函数包装为可直接实例化的节点类。两条路径共享同一个依赖注入与端口绑定基础设施(JianmuNodeMixin),这意味着无论选择哪种方式,节点都能访问 StateManager、RunContext、模型客户端、恢复 facade 以及端口 IO 能力。

类继承体系与设计原理

理解自定义节点的起点是理解 Jianmu 节点类的三层继承结构:

classDiagram
    class py_trees_behaviour_Behaviour {
        +tick() Generator
        +update() Status
        +initialise() None
        +terminate(Status) None
        +name str
        +status Status
    }
    class JianmuNodeMixin {
        +state_manager StateManager
        +ctx RunContext
        +inject(InjectPayload) None
        +bind(inputs, outputs) Self
        +read_port(name, default) Any
        +write_port(name, value) None
        +read_state(key, default) Any
        +write_state(updates) None
        +model_client ModelClient
        +interaction InteractionController
        +recovery CheckpointRecoveryFacade
    }
    class Behaviour {
        +update() Status
    }
    class AsyncBehaviour {
        +async_task Task
        +async update_async() Status
        +initialise() None
        +tick() Generator
        +update() Status
        +terminate(Status) None
    }
    class Node {
        语义别名 = Behaviour
    }
    class AsyncNode {
        语义别名 = AsyncBehaviour
    }
    class FunctionNode {
        +_func Callable
        +async update_async() Status
    }
    class StateCondition {
        +predicate Callable
        +update() Status
    }
    class Wait {
        +duration float
        +async update_async() Status
    }

    py_trees_behaviour_Behaviour <|-- Behaviour
    JianmuNodeMixin <|-- Behaviour
    Behaviour <|-- AsyncBehaviour
    Behaviour <|-- Node : 别名
    AsyncBehaviour <|-- AsyncNode : 别名
    AsyncBehaviour <|-- FunctionNode
    Behaviour <|-- StateCondition
    AsyncBehaviour <|-- Wait

最底层是 JianmuNodeMixin,它通过多重继承混入 py_trees.behaviour.Behaviour,为所有 Jianmu 节点注入了以下运行时能力:状态读写(read_state / write_state)、端口绑定与解析(read_port / write_port)、模型客户端懒加载(model_client)、人机交互控制器(interaction),以及 checkpoint-owned 恢复 facade(recovery)。Behaviour(即 Node)是同步节点的根类,你只需覆写 update() -> Status。AsyncBehaviour(即 AsyncNode)在 Behaviour 之上增加了 asyncio 任务生命周期管理——它在 initialise() 中自动创建 asyncio.Task 并注册唤醒回调,在 tick() 中检测任务完成并重启 RUNNING 态节点,在 update() 中将任务结果映射回 py_trees 状态值,在 terminate() 中取消飞行中的任务。

类继承方式:编写自定义 Node 与 AsyncNode

同步节点:继承 Node 并覆写 update()

适用于纯 CPU 计算、条件判断、简单状态变换等不涉及 IO 的场景。同步节点的 update() 方法在 py_trees 的每次 tick 中被直接调用,必须返回 Status.SUCCESS、Status.FAILURE 或 Status.RUNNING 之一。

以下是从内置实现中抽象出的最小同步节点模板:

from jianmu.node.base import Node
from py_trees.common import Status

class MySyncNode(Node):
    def __init__(self, name: str | None = None, *, namespace: str | None = None):
        resolved_name = name or self.__class__.__name__
        super().__init__(name=resolved_name, namespace=namespace)

    def update(self) -> Status:
        # 读取状态
        value = self.read_port("input_field", default=None)
        # 核心逻辑
        result = self._process(value)
        # 写回状态
        self.write_port("output_field", result)
        return Status.SUCCESS

同步节点也可用于实现控制流类型的判断节点。例如 StateCondition 就是一个典型的同步 Node 子类,它在 update() 中读取状态字段、应用比较运算符,返回 SUCCESS 或 FAILURE 来影响父级 Selector / Sequence 的路由决策:

异步节点:继承 AsyncNode 并覆写 update_async()

适用于网络调用、模型推理、文件 IO、asyncio.sleep 等异步操作。异步节点在首次进入时由 initialise() 自动将 update_async() 包装为 asyncio.Task,后续每次 tick() 都会检查任务是否完成——若完成且返回 RUNNING 则自动重新调度。

以下是从完整示例中提取的标准异步节点模板(来自 step_09_write_custom_node.py):

from jianmu.engine.behaviour import AsyncBehaviour
from py_trees.common import Status

class NormalizeTextNode(AsyncBehaviour):
    """Read one input port, normalize text, and write one output port."""

    async def update_async(self) -> Status:
        raw_text = str(self.read_port("source_text", default=""))
        normalized = " ".join(raw_text.strip().lower().split())
        self.write_port("normalized_text", normalized)
        return Status.SUCCESS

以及对应的完整编排代码:

from jianmu import Agent
from jianmu.tree import Sequence

root = Sequence(
    name="CustomNodeFlow",
    memory=True,
    children=[
        NormalizeTextNode(name="NormalizeText").bind(
            inputs={"source_text": "text"},
            outputs={"normalized_text": "normalized_text"},
        ),
        SummarizeNode(name="Summarize").bind(
            inputs={"normalized": "normalized_text"},
            outputs={"summary": "summary"},
        ),
    ],
)
agent = Agent(root, state_schema=CustomNodeState)
await agent.run()

关键要点:update_async() 必须是 async def,返回 Status 枚举值。AsyncBehaviour 基类负责管理 asyncio.Task 的创建、完成检测、取消和重启——你只需关注业务逻辑。

异步任务生命周期内幕

理解 AsyncBehaviour 的内部机制有助于调试复杂节点的行为。下表展示了节点从进入树到退出树的完整状态流转:

阶段 触发方法 内部行为 开发者关注点
进入 initialise() 取消旧任务 → 创建新 asyncio.Task(update_async()) → 注册 done_callback 唤醒 runner 在此处初始化节点级计数器、标志位
轮询 tick() 若 status==RUNNING 且 task 已完成,调用 initialise() 重新调度 无需介入
状态映射 update() 任务未完成 → 返回 RUNNING;已完成 → 返回任务结果;取消 → 返回 INVALID 无需覆写
退出 terminate() 取消飞行中的 task → 置 async_task = None 清理节点级资源

装饰器方式:用 @node 将函数包装为节点

当节点逻辑可以表达为"接收当前状态 → 返回状态更新"的纯函数时,@node 装饰器是最简洁的选择。它在内部生成一个 FunctionNode 子类(FunctionNode 本身是 AsyncBehaviour 的子类),在 update_async() 中自动完成状态读取、函数调用(同步函数在线程池中执行)、返回值校验和状态合并。

flowchart LR
    A["@node 装饰函数"] --> B["生成 WrappedNode 类"]
    B --> C["实例化: MyNode(name, state_manager)"]
    C --> D["Runner 注入: inject(payload)"]
    D --> E["initialise() 创建 asyncio.Task"]
    E --> F["update_async(): 读状态 → 调函数 → 写状态"]
    F --> G["返回 Status.SUCCESS / FAILURE"]

装饰器的两种用法:

from jianmu.node import node

# 无参数:自动使用函数名和 docstring
@node
def normalize_text(state):
    """Normalize whitespace and lowercase."""
    normalized = " ".join(state.text.strip().lower().split())
    return {"normalized": normalized}

# 带参数:显式指定名称和描述
@node(name="NormalizeText", description="Normalize text input")
def normalize_text(state):
    normalized = " ".join(state.text.strip().lower().split())
    return {"normalized": normalized}

函数签名约定:接收一个参数 state(通过 StateManager.get() 获取的 Pydantic 模型实例),返回 dict[str, Any](状态字段更新)或 None(无更新)。返回其他类型将触发 ValueError 并导致节点以 FAILURE 结束。装饰器生成的类在实例化时接受 (name, state_manager) 参数,可直接作为子节点放入 Sequence / Selector 等复合节点中:

root = Sequence(
    name="DecoratorFlow",
    memory=True,
    children=[
        normalize_text(),                           # 无参数实例化
        normalize_text(name="Step2"),               # 覆盖名称
    ],
)

如果函数是 async def,FunctionNode.update_async() 会直接 await 它;如果是同步函数,则通过 asyncio.to_thread() 在线程池中执行,避免阻塞事件循环。

运行时依赖注入

所有 Jianmu 节点(无论类继承还是装饰器)在被 ReactiveRunner 加载时都会收到一个统一的 InjectPayload,其中包含三个核心依赖:

依赖 类型 访问方式 用途
StateManager StateManager self.state_manager 读写持久化状态、订阅变更通知
RunContext RunContext self.ctx 访问模型客户端、审批管理器、沙箱、约束等共享服务
Wake-up 回调 Callable[[], None] self._wake_up 异步节点任务完成时唤醒 runner(框架自动管理)

此外,JianmuNodeMixin 还提供了三个便捷属性: - self.model_client:懒加载模型客户端,优先使用节点显式设置的 _explicit_model_client,否则回退到 ctx.model_client。 - self.interaction:返回绑定到当前 StateManager 的 InteractionController,用于触发审批、请求用户输入等人机交互流程。 - self.recovery:返回绑定到当前节点的 checkpoint-owned CheckpointRecoveryFacade。这是普通节点作者需要看到的唯一恢复入口。

依赖注入发生在 ReactiveRunner.setup() 阶段——runner 遍历整棵行为树,对每个 JianmuNodeMixin 实例调用 inject(payload)。

Checkpoint 与恢复心智

大多数自定义节点不需要写恢复代码。只要节点做的是纯计算,或者通过 read_state()、write_state()、read_port()、write_port() 读写 Jianmu state,持久状态就已经由 checkpoint 保存。

可以按这个规则判断:

  • 纯计算:不需要恢复代码。
  • 需要持久化的业务进度:写入 Jianmu state。
  • 使用内置 Tool / LLM / Agent 能力:优先使用内置节点或 API,它们会通过 checkpoint records 记录已完成的外部任务。
  • 自定义节点里直接执行外部副作用,例如 HTTP 请求、文件写入、消息发送、远端任务创建:通过 self.recovery.task(...) 和 self.recovery.tasks 登记该操作,或者显式提供幂等 / reconcile 策略。

self.recovery facade 的职责是把 task/progress records 交给 checkpoint 层管理。它不会序列化任意 Python 对象字段。不要依赖实例属性跨 checkpoint restore 存活;需要持久化的数据应从 state 或已完成的 recovery record 中重建。

dump_resume_state() 和 restore_resume_state() 是少数内置节点在旧 node-local 恢复路径迁移期间使用的内部兼容 hook,不是普通自定义节点的公开扩展模型。

实践上,大多数自定义节点根本不应该实现任何 checkpoint 专用 API。推荐模型是:

  • 持久业务进度放在普通 Jianmu state 中;
  • 只有在外部副作用需要 replay 或 reconcile 时,才通过 self.recovery 写 task records;
  • 纯计算、状态变换、幂等读取在 restore 后直接重跑即可。

这样普通自定义节点的心智会更接近 LangGraph 式的 task recovery,而不是要求每个节点作者都去设计一套自己的 dump/restore 协议。

端口绑定:声明式数据流

端口绑定是 Jianmu 节点间数据传递的核心机制。每个节点可以通过 bind() 方法声明其输入端口和输出端口与状态字段的映射关系:

MyNode(name="Processor").bind(
    inputs={"raw_input": "user_query"},     # 节点端口 "raw_input" ← 状态字段 "user_query"
    outputs={"result": "processed_output"}, # 节点端口 "result" → 状态字段 "processed_output"
)

绑定后,节点内部使用 read_port() 和 write_port() 进行读写,无需硬编码状态字段名:

class ProcessorNode(AsyncBehaviour):
    async def update_async(self) -> Status:
        # 通过端口名读取,自动解析到绑定的状态字段
        data = self.read_port("raw_input", default="")
        result = self._transform(data)
        # 通过端口名写入
        self.write_port("result", result)
        return Status.SUCCESS

这种间接层的好处是:同一个节点类可以在不同的工作流中绑定到不同的状态字段,实现真正的可复用性。如果未显式绑定,read_port("foo") 会直接以 "foo" 作为状态键查找。write_ports() 方法支持批量写入多个端口。

两种方式的对比与选择

维度 类继承(AsyncNode / Node) 装饰器(@node)
适用场景 复杂控制流、多步骤异步逻辑、需要生命周期钩子 简单状态变换、纯计算、无副作用的映射
代码量 较多:需定义类、__init__、update_async 极少:一个函数即可
状态访问 完整的 self.state_manager、self.read_port、self.write_port 通过函数参数 state 接收完整状态对象
状态写入 调用 self.write_port() 或 self.state_manager.update() 返回 dict,框架自动合并
异步支持 直接 await 任意协程 async def 函数直接支持;同步函数在线程池执行
依赖注入 完整访问 self.ctx、self.model_client、self.interaction、self.recovery 仅能通过 state_manager 间接访问(可在实例化后注入)
端口绑定 bind(inputs=..., outputs=...) 声明式绑定 默认按函数参数名推断,也可实例化后调用 .bind()
生命周期 可覆写 initialise()、terminate() 无生命周期钩子
测试 需要 StateManager 实例 同样需要 StateManager 实例

内置节点参考:学习真实实现

研究 Jianmu 内置节点的实现是掌握自定义节点编写的最佳途径。以下是值得深入阅读的实现模式:

  • 同步条件节点:StateCondition — 演示 Node.update() 中如何读取状态、应用比较运算符,返回 SUCCESS/FAILURE 来控制树路由。
  • 异步工具节点:ToolExecutor — 演示 AsyncBehaviour.update_async() 中如何调用模型客户端、处理工具调用,是复杂异步节点的标杆。
  • 轻量级异步节点:Wait — 演示最简单的 AsyncNode 用法:await asyncio.sleep() 后返回 SUCCESS。
  • 装饰器节点:Timeout — 演示 Decorator 子类如何包裹子节点、在 update() 中监控超时并取消异步任务。
  • 复合节点:LoopUntilSuccess — 演示自定义 Decorator 如何通过 self.decorated.stop(Status.INVALID) + self.state_manager.signal() 实现子节点重试循环。

继续阅读

自定义节点是 Jianmu 节点体系的基础构件。理解它之后,建议按以下路径深入: