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 节点体系的基础构件。理解它之后,建议按以下路径深入:
- 节点体系全景:AsyncNode 基类、端口绑定与依赖注入 — 深入理解端口绑定机制和依赖注入的完整流程。
- LLM 节点:AgentLLMNode 与 SimpleLLMNode 的上下文构建与模型调用 — 了解如何在内置 LLM 节点中封装模型调用与上下文构建。
- 工具与技能节点:ToolExecutor、SkillNode 与约束联动 — 了解工具执行节点如何与 Guard 体系和约束系统联动。
- 行为树执行内核:基于 py_trees 的异步扩展与 Jianmu 节点模型 — 从执行引擎视角理解节点的 tick 调度与生命周期。
- ReactiveRunner:事件驱动的异步 tick 调度与挂起恢复机制 — 理解 runner 如何注入依赖、调度 tick、处理挂起恢复。