跳转至

jianmu.node

适用对象:节点作者 / 工作流开发者 / Tree Studio 维护者 是否必读:是 相关模块:jianmu.engine, jianmu.tool, jianmu.skill, jianmu.swarm

1. 模块职责

jianmu.node 是节点公开入口,集中暴露节点基类、内置节点和预制组合。

如果说 jianmu.engine 更偏运行时核心,那么 jianmu.node 就是大多数用户真正直接编排行为树时会使用的模块。

2. 适合查什么

  • 基类与运行时节点语义:AsyncBehaviour、LoopUntilSuccess
  • 装饰器:node
  • 常见节点:AgentLLMNode、ToolExecutor、SkillNode
  • 通用节点:Log、Wait、Timeout、StateCondition
  • 多 Agent 节点:SpawnAgent、SendMessage、WaitMessage、SwarmNode
  • 预制模式:create_react_node()、create_plan_execute_node()

3. 使用建议

  • 写自定义节点,优先从这里进入
  • 使用内置工作流节点,也优先依赖这里的公开对象
  • 树组合子如 Sequence、Selector、Parallel 虽然这里仍可见,但新代码应优先从 jianmu.tree 导入
  • 如果你只是想理解“同步节点 / 异步节点”,优先看 jianmu.Node 和 jianmu.AsyncNode
  • 想理解更深层端口与注入机制,再回看 docs/concepts/node_system.md
  • create_react_node() 和 create_plan_execute_node() 现在都支持把工具集合、动态 tool provider、执行约束和 tool runner 作为 preset 入口参数传入
  • LoopUntilSuccess 和 Timeout 已统一成更偏 keyword-only 的调用风格;新代码应显式写 child=、max_iterations=、duration=
  • SpawnAgent 默认每次执行都会创建一个新 agent;如果希望复用同一个已写入状态的 agent id,可显式使用 reuse_if_present=True

4. 边界说明

  • jianmu.node 负责 jianmu 语义节点、内建节点和 preset
  • jianmu.tree 负责行为树原语
  • jianmu.node.builtin 是内建实现的组织层,不是主推荐导入路径

因此推荐:

  • from jianmu.node import Log
  • from jianmu.node import ToolExecutor

而不是默认推荐:

  • from jianmu.node.builtin import Log
  • from jianmu.node.builtin import ToolExecutor

5. 最小示例

from jianmu.model import ModelClient
from jianmu.node import ToolExecutor, create_react_node
from jianmu.tool import CalculatorTool

model_client = ModelClient.resolve()
react_node = create_react_node(
    model_client=model_client,
    tools=[CalculatorTool()],
)

tool_executor = ToolExecutor(tools=[CalculatorTool()])

6. preset 心智模型

  • create_react_node() 适合“边想边调用工具”的单循环 Agent
  • create_plan_execute_node() 适合“先规划、再执行、必要时再评审”的显式分阶段流程
  • 两者都属于公开高频 preset,文档与调用方式应尽量保持同一量级的可配置性

当前推荐的理解方式:

  • model_client:模型客户端
  • tools:静态工具集
  • tool_providers:动态工具来源
  • tool_runner:工具执行器覆盖
  • constraints:审批 / 沙箱 / 执行约束
  • config:preset 的节点级配置

7. 常见入口

  • 想写单次模型节点:看 SimpleLLMNode
  • 想让模型发起工具调用:看 AgentLLMNode + ToolExecutor
  • 想直接加载 skill:看 SkillNode
  • 想快速搭 ReAct / Plan-Execute:看 create_react_node() / create_plan_execute_node()

8. API 参考

node

jianmu Nodes: pre-built nodes for common use cases.

AsyncBehaviour

AsyncBehaviour(name: str, namespace: Optional[str] = None)

Bases: Behaviour

Base class for asynchronous Jianmu nodes.

Initialize an asynchronous Jianmu behaviour node.

源代码位于: jianmu/engine/behaviour.py
def __init__(self, name: str, namespace: Optional[str] = None):
    """Initialize an asynchronous Jianmu behaviour node."""
    super().__init__(name, namespace)
    self.async_task = None

initialise

initialise() -> None

Start a fresh async task whenever the node re-enters execution.

源代码位于: jianmu/engine/behaviour.py
def initialise(self) -> None:
    """Start a fresh async task whenever the node re-enters execution."""
    if self.async_task and not self.async_task.done():
        self.async_task.cancel()

    try:
        loop = asyncio.get_running_loop()
        self.async_task = loop.create_task(self.update_async())

        # Waking the runner on completion keeps event-driven execution responsive
        # without requiring the runner to poll for task state.
        wake_up = self._wake_up
        if wake_up:
            self.async_task.add_done_callback(lambda _: wake_up())
    except RuntimeError:
        self.feedback_message = "❌ No active asyncio event loop found."
        self.async_task = None

tick

tick()

Restart completed RUNNING tasks so async nodes can make forward progress across ticks.

源代码位于: jianmu/engine/behaviour.py
def tick(self):
    """Restart completed RUNNING tasks so async nodes can make forward progress across ticks."""
    if self.status == Status.RUNNING and self.async_task is not None and self.async_task.done():
        try:
            previous_status = self.async_task.result()
        except asyncio.CancelledError:
            previous_status = Status.INVALID
        except Exception:
            previous_status = Status.FAILURE
        if previous_status == Status.RUNNING:
            self.initialise()
    for node in super().tick():
        yield node

update

update() -> Status

Map async task state back into py_trees status values.

返回:

类型 描述
Status

RUNNING while the async task is in progress, FAILURE if the

Status

task is missing or returned an invalid type, INVALID if it was

Status

cancelled, or the Status value returned by update_async.

源代码位于: jianmu/engine/behaviour.py
def update(self) -> Status:
    """Map async task state back into py_trees status values.

    Returns:
        ``RUNNING`` while the async task is in progress, ``FAILURE`` if the
        task is missing or returned an invalid type, ``INVALID`` if it was
        cancelled, or the ``Status`` value returned by ``update_async``.
    """
    if self.async_task is None:
        return Status.FAILURE

    if not self.async_task.done():
        return Status.RUNNING

    try:
        status = self.async_task.result()
        if not isinstance(status, Status):
            self.feedback_message = f"Invalid return type: {type(status)}"
            return Status.FAILURE
        return status

    except asyncio.CancelledError:
        return Status.INVALID

terminate

terminate(new_status: Status) -> None

Cancel the in-flight task when the node is interrupted.

参数:

名称 类型 描述 默认
new_status Status

The status py_trees is transitioning this node to.

必需
源代码位于: jianmu/engine/behaviour.py
def terminate(self, new_status: Status) -> None:
    """Cancel the in-flight task when the node is interrupted.

    Args:
        new_status: The status py_trees is transitioning this node to.
    """
    if self.async_task and not self.async_task.done():
        self.async_task.cancel()
    self.async_task = None

update_async async

update_async() -> Status

Run one asynchronous node update.

返回:

类型 描述
Status

py_trees status value produced by the async node logic.

引发:

类型 描述
NotImplementedError

Always raised in the base class; subclasses must override this method.

源代码位于: jianmu/engine/behaviour.py
async def update_async(self) -> Status:
    """Run one asynchronous node update.

    Returns:
        py_trees status value produced by the async node logic.

    Raises:
        NotImplementedError: Always raised in the base class; subclasses
            must override this method.
    """
    raise NotImplementedError("AsyncBehaviour subclass must implement update_async()")

LoopUntilSuccess

LoopUntilSuccess(
    name: str | None = None,
    *,
    child: Optional[Behaviour] = None,
    max_iterations: int = 10,
    abort_condition: Optional[Callable[[], bool]] = None,
)

Bases: Decorator

Retry a child until it succeeds or the iteration budget is exhausted.

This decorator converts child failure into another execution round until the child succeeds, the retry budget is exhausted, or an abort condition triggers.

属性:

名称 类型 描述
max_iterations

Maximum allowed retry count before terminal failure.

iteration_count

Current retry count for the active entry.

abort_condition

Optional callable used to stop the loop early.

Configure retry-until-success behaviour with an abort hook.

参数:

名称 类型 描述 默认
name str | None

Behavior tree node name.

None
child Optional[Behaviour]

Wrapped child behavior to retry.

None
max_iterations int

Maximum number of agent-loop rounds triggered by child failure before the loop fails terminally.

10
abort_condition Optional[Callable[[], bool]]

Optional callable that stops the loop early when it returns True.

None
源代码位于: jianmu/node/composites.py
def __init__(
    self,
    name: str | None = None,
    *,
    child: Optional[Behaviour] = None,
    max_iterations: int = 10,
    abort_condition: Optional[Callable[[], bool]] = None,
):
    """Configure retry-until-success behaviour with an abort hook.

    Args:
        name: Behavior tree node name.
        child: Wrapped child behavior to retry.
        max_iterations: Maximum number of agent-loop rounds triggered by
            child failure before the loop fails terminally.
        abort_condition: Optional callable that stops the loop early when
            it returns ``True``.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, child=child)
    self.max_iterations = max_iterations
    self.iteration_count = 0
    self.abort_condition = abort_condition

dump_resume_state

dump_resume_state() -> dict

Return loop-local state needed to continue after checkpoint restore.

源代码位于: jianmu/node/composites.py
def dump_resume_state(self) -> dict:
    """Return loop-local state needed to continue after checkpoint restore."""
    self._write_iteration_progress()
    return {"iteration_count": int(self.iteration_count)}

restore_resume_state

restore_resume_state(data: dict) -> None

Restore loop-local retry progress from checkpoint metadata.

源代码位于: jianmu/node/composites.py
def restore_resume_state(self, data: dict) -> None:
    """Restore loop-local retry progress from checkpoint metadata."""
    try:
        self.iteration_count = max(0, int(dict(data or {}).get("iteration_count") or 0))
    except Exception:
        self.iteration_count = 0

initialise

initialise() -> None

Reset the retry counter at the start of each entry.

源代码位于: jianmu/node/composites.py
def initialise(self) -> None:
    """Reset the retry counter at the start of each entry."""
    restored = self.consume_restored_entry() == "RUNNING"
    if restored:
        progress_count = self._read_iteration_progress()
        if progress_count is not None:
            self.iteration_count = progress_count
    else:
        self.iteration_count = 0
        self._write_iteration_progress()
    if self.state_manager is not None:
        self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)

update

update() -> Status

Translate child failure into a scheduled retry.

返回:

类型 描述
Status

The current decorator status after evaluating the child state.

源代码位于: jianmu/node/composites.py
def update(self) -> Status:
    """Translate child failure into a scheduled retry.

    Returns:
        The current decorator status after evaluating the child state.
    """
    if not self.decorated:
        return Status.FAILURE

    if self.abort_condition and self.abort_condition():
        logger.error("🛑 [{}] Hard abort condition met (for example: budget exceeded); aborting loop", self.name)
        return Status.FAILURE

    child_status = self.decorated.status

    if child_status == Status.SUCCESS:
        logger.debug("✅ [{}] Loop completed successfully ({} iterations)", self.name, self.iteration_count)
        return Status.SUCCESS

    if child_status == Status.RUNNING:
        return Status.RUNNING

    if child_status == Status.FAILURE:
        if self.state_manager is not None:
            termination = self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
            terminal_reason = str(termination.get("reason") or "") if isinstance(termination, dict) else ""
            if terminal_reason in {"provider_error", "provider_retry_exhausted", "model_call_failed", "hook_error"}:
                logger.warning(
                    "⚠️ [{}] Terminal model path error detected ({}); stopping loop without retry",
                    self.name,
                    terminal_reason,
                )
                return Status.FAILURE
        self.iteration_count += 1
        self._write_iteration_progress()
        if self.iteration_count >= self.max_iterations:
            logger.warning("⚠️ [{}] Reached max iterations ({}); forcing stop", self.name, self.max_iterations)
            self._record_max_iterations_termination()
            return Status.FAILURE

        logger.debug(
            "🔄 [{}] Iteration {} incomplete; continuing to next round (max={})",
            self.name,
            self.iteration_count,
            self.max_iterations,
        )

        self.decorated.stop(Status.INVALID)

        if self.state_manager is not None:
            try:
                loop = asyncio.get_running_loop()
                loop.call_soon(self.state_manager.signal)
            except RuntimeError:
                self.state_manager.signal()

        return Status.RUNNING

    return Status.INVALID

terminate

terminate(new_status: Status) -> None

Clear retry bookkeeping on exit.

参数:

名称 类型 描述 默认
new_status Status

Final status assigned to this decorator.

必需
源代码位于: jianmu/node/composites.py
def terminate(self, new_status: Status) -> None:
    """Clear retry bookkeeping on exit.

    Args:
        new_status: Final status assigned to this decorator.
    """
    self.iteration_count = 0
    if self.decorated and self.decorated.status != Status.INVALID:
        self.decorated.stop(new_status)

AgentLLMNode

AgentLLMNode(
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    tools_schema: Optional[List[Dict[str, Any]]] = None,
    tools_description: str = "",
    context_builder: Optional[
        ContextBuilderProtocol
    ] = None,
    config: Optional[AgentLLMConfig] = None,
)

Bases: SimpleLLMNode

Run an agent-oriented LLM step that may emit structured tool calls.

Compared with SimpleLLMNode, this node additionally exposes tool schemas to the model, accumulates usage across iterations, and persists emitted tool calls into shared state for ToolExecutor to consume.

属性:

名称 类型 描述
_tools_schema

Structured tool schemas exposed to the model.

_tools_description

Human-readable tool description block for prompts.

Initialize the node.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace Optional[str]

Optional state namespace for port resolution.

None
model_client ModelClient

Explicit model client used for inference.

必需
tools_schema Optional[List[Dict[str, Any]]]

Structured tool schema exposed to the model.

None
tools_description str

Human-readable tool description block for prompts.

''
context_builder Optional[ContextBuilderProtocol]

Optional context builder override.

None
config Optional[AgentLLMConfig]

Agent LLM execution configuration.

None
源代码位于: jianmu/node/builtin/llm.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    tools_schema: Optional[List[Dict[str, Any]]] = None,
    tools_description: str = "",
    context_builder: Optional[ContextBuilderProtocol] = None,
    config: Optional[AgentLLMConfig] = None,
):
    """Initialize the node.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        model_client: Explicit model client used for inference.
        tools_schema: Structured tool schema exposed to the model.
        tools_description: Human-readable tool description block for prompts.
        context_builder: Optional context builder override.
        config: Agent LLM execution configuration.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(
        name=resolved_name,
        model_client=model_client,
        context_builder=context_builder,
        namespace=namespace,
        config=config or AgentLLMConfig(),
    )

    self._tools_schema = tools_schema
    self._tools_description = tools_description

update_async async

update_async() -> Status

Run one agent step and extract any emitted tool calls.

返回:

类型 描述
Status

Status.SUCCESS when the model step completes, otherwise

Status

Status.FAILURE.

源代码位于: jianmu/node/builtin/llm.py
async def update_async(self) -> Status:
    """Run one agent step and extract any emitted tool calls.

    Returns:
        ``Status.SUCCESS`` when the model step completes, otherwise
        ``Status.FAILURE``.
    """
    try:
        tools_schema = list(self._tools_schema or [])
        if self.state_manager is not None:
            self.state_manager.reset_call_fields()
        model_client = self.model_client or ModelClient.resolve()
        model_client = require_model_client(
            model_client,
            source=f"{self.__class__.__name__}.model_client",
        )
        global_state = self.state_manager.get() if self.state_manager else {}
        messages = self._prepare_messages()

        if not messages:
            logger.warning("⚠️ [{}] No messages, cannot call LLM", self.name)
            return Status.FAILURE

        full_messages = self._resolve_context_builder().build(
            local_state={"messages": messages},
            global_state=global_state,
            ctx=self.ctx,
            tools_schema=tools_schema,
        )
        model_config = self._build_model_config(tools_schema=tools_schema)
        full_messages, model_config = await self._apply_before_model_hooks(
            messages=full_messages,
            config=model_config,
            tools_schema=tools_schema,
            state_snapshot=global_state,
        )
        task = self._llm_task_record(
            messages=full_messages,
            config=model_config,
            tools_schema=tools_schema,
            mode="agent",
        )
        completed_record = self._completed_llm_record(
            str(task.get("task_id") or ""),
            input_hash=str(task.get("input_hash") or ""),
        )
        if completed_record is not None:
            response_msg = self._message_from_record(completed_record)
            if response_msg is not None:
                content = str(completed_record.get("content") or "")
                self._emit_replayed_llm_events(
                    response_msg=response_msg,
                    content=content,
                    model_config=model_config,
                )
                self._persist_usage(response_msg)
                self._persist_result(response_msg=response_msg, content=content)
                return Status.SUCCESS
        self._emit_runtime_event(
            "reply.started",
            payload={
                "message_count": len(messages),
                "full_message_count": len(full_messages),
                "stream": self.config.stream,
            },
        )
        self._emit_runtime_event(
            "llm.request",
            payload={
                "model": model_config.model,
                "stream": self.config.stream,
                "message_count": len(full_messages),
                "messages": [
                    {
                        key: value
                        for key, value in message.to_dict().items()
                        if key not in {"id", "tool"}
                    }
                    for message in full_messages
                ],
            },
        )
        self._emit_runtime_event(
            "model.call.started",
            payload={
                "model": model_config.model,
                "stream": self.config.stream,
                "tool_schema_count": len(model_config.tools or []),
            },
        )
        on_text_update = None
        if self.config.keys.streaming_output or self.config.stream:
            delta_emitter = (
                TextDeltaEmitter(lambda event_type, payload: self._emit_runtime_event(event_type, payload=payload))
                if self.config.stream
                else None
            )

            def on_text_update(current_text: str) -> None:
                if self.config.keys.streaming_output:
                    self.write_port(
                        "streaming_output",
                        current_text,
                        signal=False,
                        state_key=self.config.keys.streaming_output,
                    )
                if delta_emitter is not None:
                    delta_emitter.push(current_text)

        self.recovery.tasks.begin(task)
        with span("llm_call", model=model_config.model):
            response_msg, content = await model_client.invoke(
                messages=full_messages,
                config=model_config,
                stream=self.config.stream,
                on_text_update=on_text_update,
                trace_node=self.name,
            )
        response_msg, content = await self._apply_after_model_hooks(
            response_msg=response_msg,
            content=content,
            tools_schema=tools_schema,
            state_snapshot=global_state,
        )
        usage = dict((response_msg.metadata or {}).get("usage") or {})
        tool_call_count = len(getattr(response_msg, "tool_calls", []) or [])
        self._emit_runtime_event(
            "model.call.completed",
            payload={
                "model": model_config.model,
                "tool_call_count": tool_call_count,
                "usage": usage,
            },
        )
        self._emit_runtime_event(
            "reply.completed",
            payload={
                "final_text": content,
                "tool_call_count": tool_call_count,
            },
        )

        self._persist_usage(response_msg)
        self._persist_result(response_msg=response_msg, content=content)
        self._write_completed_llm_record(
            task=task,
            response_msg=response_msg,
            content=content,
            tools_schema=tools_schema,
            mode="agent",
        )
        return Status.SUCCESS
    except Exception as e:
        await self._notify_model_error_hooks(error=e, tools_schema=locals().get("tools_schema") or [])
        self._record_model_failure(e)
        logger.error("🔥 [{}] AgentLLM failed: {}", self.name, e)
        return Status.FAILURE

SimpleLLMNode

SimpleLLMNode(
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    context_builder: Optional[
        ContextBuilderProtocol
    ] = None,
    config: Optional[AgentLLMConfig] = None,
)

Bases: AsyncNode

Run one LLM call and persist the response back to state.

SimpleLLMNode is the minimal model-inference node: it builds prompt context from state, calls a model client once, appends the assistant message, and writes the resulting text to text_output.

Use this when you need one model turn without tool calling. For tool-capable agent turns, use AgentLLMNode.

属性:

名称 类型 描述
_explicit_model_client

Explicit model client bound at construction time.

config

LLM execution configuration for the node.

_explicit_context_builder

Optional explicit context-builder override.

Initialize the node.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace Optional[str]

Optional state namespace for port resolution.

None
model_client ModelClient

Explicit model client used for inference.

必需
context_builder Optional[ContextBuilderProtocol]

Optional context builder override. When omitted, the node resolves a default chat or ReAct builder.

None
config Optional[AgentLLMConfig]

LLM execution configuration.

None
源代码位于: jianmu/node/builtin/llm.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    context_builder: Optional[ContextBuilderProtocol] = None,
    config: Optional[AgentLLMConfig] = None,
):
    """Initialize the node.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        model_client: Explicit model client used for inference.
        context_builder: Optional context builder override. When omitted,
            the node resolves a default chat or ReAct builder.
        config: LLM execution configuration.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(resolved_name, namespace=namespace)
    self._explicit_model_client = require_model_client(
        model_client,
        source=f"{self.__class__.__name__}.model_client",
    )
    self.config = config or AgentLLMConfig()
    self._explicit_context_builder = context_builder

initialise

initialise() -> None

Start an async model call and reset stale task records on fresh entry.

源代码位于: jianmu/node/builtin/llm.py
def initialise(self) -> None:
    """Start an async model call and reset stale task records on fresh entry."""
    if self.consume_restored_entry() != "RUNNING":
        self.recovery.tasks.clear_owner()
    super().initialise()

update_async async

update_async() -> Status

Build prompt context, invoke the model, and persist the answer.

返回:

类型 描述
Status

Status.SUCCESS when a response is written successfully, otherwise

Status

Status.FAILURE.

源代码位于: jianmu/node/builtin/llm.py
async def update_async(self) -> Status:
    """Build prompt context, invoke the model, and persist the answer.

    Returns:
        ``Status.SUCCESS`` when a response is written successfully, otherwise
        ``Status.FAILURE``.
    """
    try:
        if self.state_manager is not None:
            self.state_manager.reset_call_fields()
        messages = self._prepare_messages()
        if not messages:
            logger.warning("⚠️ [{}] No messages, cannot call LLM", self.name)
            return Status.FAILURE

        model_client = self.model_client or ModelClient.resolve()
        model_client = require_model_client(
            model_client,
            source=f"{self.__class__.__name__}.model_client",
        )
        global_state = self.state_manager.get() if self.state_manager else {}
        full_messages = self._resolve_context_builder().build(
            local_state={"messages": messages},
            global_state=global_state,
            ctx=self.ctx,
            tools_schema=[],
        )
        model_config = self._build_model_config()
        full_messages, model_config = await self._apply_before_model_hooks(
            messages=full_messages,
            config=model_config,
            tools_schema=[],
            state_snapshot=global_state,
        )
        task = self._llm_task_record(
            messages=full_messages,
            config=model_config,
            tools_schema=[],
            mode="simple",
        )
        completed_record = self._completed_llm_record(
            str(task.get("task_id") or ""),
            input_hash=str(task.get("input_hash") or ""),
        )
        if completed_record is not None:
            response_msg = self._message_from_record(completed_record)
            if response_msg is not None:
                content = str(completed_record.get("content") or "")
                self._emit_replayed_llm_events(
                    response_msg=response_msg,
                    content=content,
                    model_config=model_config,
                )
                self.append_port_messages(
                    "messages",
                    [response_msg],
                    state_key=self.config.keys.messages,
                    signal=False,
                )
                self._persist_text_output(content)
                return Status.SUCCESS
        self._emit_runtime_event(
            "reply.started",
            payload={
                "message_count": len(messages),
                "full_message_count": len(full_messages),
                "stream": self.config.stream,
            },
        )
        self._emit_runtime_event(
            "llm.request",
            payload={
                "model": model_config.model,
                "stream": self.config.stream,
                "message_count": len(full_messages),
                "messages": [
                    {
                        key: value
                        for key, value in message.to_dict().items()
                        if key not in {"id", "tool"}
                    }
                    for message in full_messages
                ],
            },
        )
        self._emit_runtime_event(
            "model.call.started",
            payload={
                "model": model_config.model,
                "stream": self.config.stream,
                "tool_schema_count": len(model_config.tools or []),
            },
        )
        on_text_update = None
        if self.config.keys.streaming_output or self.config.stream:
            delta_emitter = (
                TextDeltaEmitter(lambda event_type, payload: self._emit_runtime_event(event_type, payload=payload))
                if self.config.stream
                else None
            )

            def on_text_update(current_text: str) -> None:
                if self.config.keys.streaming_output:
                    self.write_port(
                        "streaming_output",
                        current_text,
                        signal=False,
                        state_key=self.config.keys.streaming_output,
                    )
                if delta_emitter is not None:
                    delta_emitter.push(current_text)

        self.recovery.tasks.begin(task)
        with span("llm_call", model=model_config.model):
            response_msg, content = await model_client.invoke(
                messages=full_messages,
                config=model_config,
                stream=self.config.stream,
                on_text_update=on_text_update,
                trace_node=self.name,
            )
        response_msg, content = await self._apply_after_model_hooks(
            response_msg=response_msg,
            content=content,
            tools_schema=[],
            state_snapshot=global_state,
        )
        usage = dict((response_msg.metadata or {}).get("usage") or {})
        tool_call_count = len(getattr(response_msg, "tool_calls", []) or [])
        self._emit_runtime_event(
            "model.call.completed",
            payload={
                "model": model_config.model,
                "tool_call_count": tool_call_count,
                "usage": usage,
            },
        )
        self._emit_runtime_event(
            "reply.completed",
            payload={
                "final_text": content,
                "tool_call_count": tool_call_count,
            },
        )

        self.append_port_messages(
            "messages",
            [response_msg],
            state_key=self.config.keys.messages,
            signal=False,
        )
        self._persist_text_output(content)
        self._write_completed_llm_record(
            task=task,
            response_msg=response_msg,
            content=content,
            tools_schema=[],
            mode="simple",
        )
        return Status.SUCCESS
    except Exception as e:
        await self._notify_model_error_hooks(error=e, tools_schema=[])
        self._record_model_failure(e)
        logger.error("🔥 [{}] LLM call failed: {}", self.name, e)
        return Status.FAILURE

ToolExecutor

ToolExecutor(
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    tools: Optional[List[Tool]] = None,
    tool_runner: Optional[Any] = None,
    guard_enforcer: Optional[Any] = None,
    constraints: Optional[Any] = None,
    config: Optional[ToolExecutorConfig] = None,
    tool_result_policy: Optional[ToolResultPolicy] = None,
    tool_result_policy_scope: Optional[
        ToolResultPolicyScope
    ] = None,
)

Bases: AsyncNode

Execute agent-selected tool calls and append observations to state.

ToolExecutor expects pending tool calls to already exist in agent state and is typically paired with AgentLLMNode or one of the preset agent loops. It resolves the requested tool by name, runs it, and appends a tool observation message back into the shared conversation state.

属性:

名称 类型 描述
toolset

Registered runtime tool collection visible to the executor.

config

Execution and observation configuration for tool runs.

constraints

Optional guard and execution constraints for the batch.

Initialize the executor.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace Optional[str]

Optional state namespace for port resolution.

None
tools Optional[List[Tool]]

Tool registry exposed to the agent.

None
tool_runner Optional[Any]

Optional explicit tool runner override.

None
guard_enforcer Optional[Any]

Optional explicit guard enforcer override.

None
constraints Optional[Any]

Optional guard and execution constraints used to build approval and execution policy components.

None
config Optional[ToolExecutorConfig]

Executor configuration including retry and key settings.

None
tool_result_policy Optional[ToolResultPolicy]

Optional policy evaluated after each completed tool call and before its observation is committed.

None
tool_result_policy_scope Optional[ToolResultPolicyScope]

Optional predicate limiting which tool calls are handled by tool_result_policy.

None
源代码位于: jianmu/node/builtin/tool.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    tools: Optional[List[Tool]] = None,
    tool_runner: Optional[Any] = None,
    guard_enforcer: Optional[Any] = None,
    constraints: Optional[Any] = None,
    config: Optional[ToolExecutorConfig] = None,
    tool_result_policy: Optional[ToolResultPolicy] = None,
    tool_result_policy_scope: Optional[ToolResultPolicyScope] = None,
):
    """Initialize the executor.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        tools: Tool registry exposed to the agent.
        tool_runner: Optional explicit tool runner override.
        guard_enforcer: Optional explicit guard enforcer override.
        constraints: Optional guard and execution constraints used to build
            approval and execution policy components.
        config: Executor configuration including retry and key settings.
        tool_result_policy: Optional policy evaluated after each completed
            tool call and before its observation is committed.
        tool_result_policy_scope: Optional predicate limiting which tool
            calls are handled by ``tool_result_policy``.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.toolset = ToolSet()
    self.config = config or ToolExecutorConfig()

    self._tool_runner_override = tool_runner
    self._tool_runner = None
    self._enforcer_override = guard_enforcer
    self._enforcer = guard_enforcer
    self.constraints = constraints
    self.tool_result_policy = tool_result_policy
    self.tool_result_policy_scope = tool_result_policy_scope
    self._completion_status: Status | None = None
    self._inflight_tool_tasks: set[asyncio.Task] = set()
    self._pending_enforcer_resume_state: dict | None = None

    if tools:
        for t in tools:
            self.register_tool(t)

completion_status property

completion_status: Status | None

Return the status requested by the latest Complete disposition.

返回:

类型 描述
Status | None

Completion status for the active batch, or None when no policy

Status | None

result has completed the ReAct loop.

initialise

initialise() -> None

Reset policy completion state and start a fresh asynchronous batch.

源代码位于: jianmu/node/builtin/tool.py
def initialise(self) -> None:
    """Reset policy completion state and start a fresh asynchronous batch."""
    if self.consume_restored_entry() != "RUNNING":
        self.recovery.tasks.clear_owner()
    self._completion_status = None
    super().initialise()

dump_resume_state

dump_resume_state() -> dict

Return tool-executor local state needed for checkpoint resume.

源代码位于: jianmu/node/builtin/tool.py
def dump_resume_state(self) -> dict:
    """Return tool-executor local state needed for checkpoint resume."""
    if self._enforcer is None:
        return {}
    dump_resume_state = getattr(self._enforcer, "dump_resume_state", None)
    if not callable(dump_resume_state):
        return {}
    enforcer_state = dump_resume_state()
    return {"guard_enforcer": enforcer_state} if isinstance(enforcer_state, dict) and enforcer_state else {}

restore_resume_state

restore_resume_state(data: dict) -> None

Restore pending tool-executor state from checkpoint metadata.

源代码位于: jianmu/node/builtin/tool.py
def restore_resume_state(self, data: dict) -> None:
    """Restore pending tool-executor state from checkpoint metadata."""
    guard_state = dict(data or {}).get("guard_enforcer")
    self._pending_enforcer_resume_state = dict(guard_state) if isinstance(guard_state, dict) else None

inject

inject(payload: Any) -> None

Bind runtime dependencies and cancel workers when a run detaches.

源代码位于: jianmu/node/builtin/tool.py
def inject(self, payload: Any) -> None:
    """Bind runtime dependencies and cancel workers when a run detaches."""
    detaching = getattr(self, "_wake_up", None) is not None and getattr(payload, "wake_up", None) is None
    super().inject(payload)
    if detaching:
        self._cancel_inflight_tool_tasks()
        if self.async_task is not None and not self.async_task.done():
            self.async_task.cancel()

terminate

terminate(new_status: Status) -> None

Cancel child tool workers before interrupting the coordinator task.

源代码位于: jianmu/node/builtin/tool.py
def terminate(self, new_status: Status) -> None:
    """Cancel child tool workers before interrupting the coordinator task."""
    self._cancel_inflight_tool_tasks()
    super().terminate(new_status)

register_tool

register_tool(tool: Tool)

Register a tool under its lowercase runtime name.

This is a supported extension point for workflow assembly and visual editors that attach ToolNode definitions to a ToolExecutor.

参数:

名称 类型 描述 默认
tool Tool

Tool instance to expose for execution.

必需
源代码位于: jianmu/node/builtin/tool.py
def register_tool(self, tool: Tool):
    """Register a tool under its lowercase runtime name.

    This is a supported extension point for workflow assembly and visual
    editors that attach ``ToolNode`` definitions to a ``ToolExecutor``.

    Args:
        tool: Tool instance to expose for execution.
    """
    self.toolset = self.toolset.with_tools([tool])
    logger.debug("🔧 [{}] Registered tool: {}", self.name, tool.name)

update_async async

update_async() -> Status

Execute the current batch of agent-requested tool calls.

返回:

类型 描述
Status

Status.SUCCESS when execution completes or no actions are pending.

Status

Status.FAILURE when state is unavailable for observation writes.

源代码位于: jianmu/node/builtin/tool.py
async def update_async(self) -> Status:
    """Execute the current batch of agent-requested tool calls.

    Returns:
        ``Status.SUCCESS`` when execution completes or no actions are pending.
        ``Status.FAILURE`` when state is unavailable for observation writes.
    """
    # Resolve guard components once per run configuration.
    if self._enforcer is None or self._tool_runner is None:
        from jianmu.guard.enforcer import GuardEnforcer
        from jianmu.execution.runner import ToolRunner

        constraints = self.constraints or (self.ctx.constraints if self.ctx else None)
        guard_constraints = constraints.guard if constraints else None
        execution_constraints = constraints.execution if constraints else None

        if self._enforcer_override is not None:
            self._enforcer = self._enforcer_override
        else:
            self._enforcer = GuardEnforcer.from_guard_constraints(
                guard_constraints,
                approval_manager=self.ctx.approval_manager if self.ctx else None,
            )
        self._tool_runner = ToolRunner.from_execution_constraints(
            execution=execution_constraints,
            builtin_names=self.toolset.names(),
            runtime_event_bus=getattr(self.ctx, "runtime_event_bus", None) if self.ctx is not None else None,
        )
        if self._tool_runner is not None and getattr(self.config, "timeout_s", None) is not None:
            self._tool_runner._timeout_s = float(self.config.timeout_s)
        # Allow callers to force a specific runner implementation.
        if self._tool_runner_override:
            self._tool_runner = self._tool_runner_override
        self._apply_pending_enforcer_resume_state()

    # Read pre-parsed actions produced by ``AgentLLMNode``.
    raw_actions = self.read_port("actions", [], state_key=self.config.keys.actions)
    actions = self._ensure_action_ids(ToolCall.from_list(raw_actions))
    actions = self._suspended_tool_store().restore_actions(actions)

    if not actions:
        return Status.SUCCESS
    try:
        actions = await self._apply_before_tool_hooks_to_actions(actions)
    except Exception as exc:
        write_runtime_termination(
            self.state_manager,
            build_hook_error_termination(exc, node_name=self.name),
        )
        logger.warning("⚠️ [{}] before_tool_call hook failed: {}", self.name, exc)
        return Status.FAILURE

    batches = self._plan_action_batches(actions)
    replay_count, replayed_results, replay_completion = self._completed_replay_prefix(actions, batches)
    actions_to_execute = actions[replay_count:]
    replayed_actions = actions[:replay_count]

    if replay_completion is not None and actions_to_execute:
        actions_to_execute = []

    decision, reason, approved_request_ids = await self._preflight_actions(actions_to_execute)
    if decision == Decision.PENDING:
        return Status.RUNNING
    if decision in (Decision.DENY, Decision.CONFIRM):
        self._record_tool_deny_termination(
            action=actions_to_execute[0] if actions_to_execute else None,
            decision=decision,
            reason=reason,
        )
        observations = [
            self._build_observation(result, tool_call_id=action.id)
            for action, result in zip(replayed_actions, replayed_results)
        ] + [
            self._build_observation(
                ToolResult.from_raw(action.name, None, error=reason or "Denied by policy"),
                tool_call_id=action.id,
            )
            for action in actions_to_execute
        ]
        self.append_port_messages(
            "messages",
            observations,
            state_key=self.config.keys.messages,
        )
        return Status.SUCCESS

    # Execute tools either serially or in parallel based on policy.
    approval = self._approval_runtime_for_batch()
    for action in actions:
        action_id = str(action.id or "").strip()
        if action_id in approved_request_ids:
            approval.write_progress(
                action.name,
                action_id,
                state=ApprovalProgressState.APPROVED.value,
                request_id=approved_request_ids[action_id],
            )
    batches = self._plan_action_batches(actions)
    self._emit_action_batch_plan(batches)
    suspended_action: ToolCall | None = None
    completion: Complete | None = replay_completion
    executed_results: list[ToolResult] = []
    try:
        if actions_to_execute and completion is None:
            execute_batches = self._plan_action_batches(actions_to_execute)
            executed_results, suspended_action, completion = await self._execute_action_batches(execute_batches)
    except Exception as exc:
        write_runtime_termination(
            self.state_manager,
            build_hook_error_termination(exc, node_name=self.name),
        )
        logger.warning("⚠️ [{}] tool result handling failed: {}", self.name, exc)
        return Status.FAILURE

    if self._suspended_tool_store().persist_active(
        actions=actions,
        suspended_action=suspended_action,
    ):
        return Status.RUNNING

    # Commit observations only after the optional result policy has run.
    if self.state_manager is None:
        return Status.FAILURE
    results = [*replayed_results, *executed_results]
    committed_count = len(results)
    if completion is not None:
        committed_count -= 1
        self._rewrite_pending_assistant_tool_calls(keep_count=committed_count)
    observations = [
        self._build_observation(result, tool_call_id=action.id)
        for action, result in zip(
            actions[:committed_count],
            results[:committed_count],
        )
    ]
    if observations:
        self.append_port_messages(
            "messages",
            observations,
            state_key=self.config.keys.messages,
        )
    self._suspended_tool_store().clear()
    for action in actions:
        approval.clear_progress(action.id)
    self._write_tool_effects(actions[: len(results)], results)
    if completion is not None:
        updates = dict(completion.state_patch)
        updates[self.config.keys.actions] = []
        self.write_state(updates, signal=False)
        self._completion_status = completion.status
        return completion.status
    return Status.SUCCESS

reset_for_run

reset_for_run() -> None

Reset guard state between agent runs.

源代码位于: jianmu/node/builtin/tool.py
def reset_for_run(self) -> None:
    """Reset guard state between agent runs."""
    if self._enforcer:
        self._enforcer.reset()
        self._apply_pending_enforcer_resume_state()

ToolNode

ToolNode(
    name: str | None = None,
    *,
    namespace: str | None = None,
    tool: Tool,
    input_key: str = "input",
    output_key: str = "output",
    execute: Optional[bool] = None,
)

Bases: AsyncNode

Wrap one Tool so it can participate directly in a behavior tree.

ToolNode is the right abstraction when a workflow explicitly decides where a tool call sits in the tree. In contrast, ToolExecutor is used by agent loops that receive tool calls from an LLM.

属性:

名称 类型 描述
tool

Underlying tool implementation invoked by the node.

input_key

State key used to read tool input.

output_key

State key used to persist tool output.

execute

Whether the tool should execute when the node is ticked.

_tool_runner

Cached sandbox-aware tool runner instance.

_tool_runner_signature

Stable signature used to decide runner reuse.

Initialize the node.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
tool Tool

Underlying tool implementation to invoke.

必需
input_key str

State key used to read tool input.

'input'
output_key str

State key used to persist tool output.

'output'
execute Optional[bool]

Whether the node should execute the tool when ticked.

None
源代码位于: jianmu/node/builtin/tool.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    tool: Tool,
    input_key: str = "input",
    output_key: str = "output",
    execute: Optional[bool] = None,
):
    """Initialize the node.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        tool: Underlying tool implementation to invoke.
        input_key: State key used to read tool input.
        output_key: State key used to persist tool output.
        execute: Whether the node should execute the tool when ticked.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.tool = tool
    self.input_key = input_key
    self.output_key = output_key
    self.execute = True if execute is None else bool(execute)
    self._tool_runner = None
    self._tool_runner_signature = None

update_async async

update_async() -> Status

Run the wrapped tool once and write its output back to state.

返回:

类型 描述
Status

Status.SUCCESS when the tool completes successfully, otherwise

Status

Status.FAILURE.

源代码位于: jianmu/node/builtin/tool.py
async def update_async(self) -> Status:
    """Run the wrapped tool once and write its output back to state.

    Returns:
        ``Status.SUCCESS`` when the tool completes successfully, otherwise
        ``Status.FAILURE``.
    """
    if not self.state_manager:
        logger.error("❌ [{}] state_manager was not injected", self.name)
        return Status.FAILURE

    if not self.execute:
        return Status.SUCCESS

    try:
        tool_args = self._resolve_inputs(allow_non_dict=True)
        tool_name = getattr(self.tool, "name", type(self.tool).__name__)
        tool_args = await _apply_before_tool_hooks(
            self,
            tool_name=tool_name,
            tool_args=tool_args,
            mode="workflow",
            tool_obj=self.tool,
        )
        _emit_node_runtime_event(
            self,
            "tool.call.started",
            payload={
                "tool": tool_name,
                "mode": "workflow",
                "args_preview": self._trace_preview_payload(tool_args),
            },
        )
        trace_emit("tool.call", {
            "node": self.name,
            "tool": tool_name,
            "mode": "workflow",
            "args": tool_args,
        })
        self._emit_bt_skill_node_io(
            stage="started",
            ok=None,
            tool_name=tool_name,
            input_value=tool_args,
            output_value=None,
            error=None,
        )
        result = await self._run_tool(tool_args)
        result = await _apply_after_tool_hooks(
            self,
            tool_name=tool_name,
            tool_args=tool_args,
            result=result,
            mode="workflow",
            tool_obj=self.tool,
        )
        tool_result = ToolResult.from_raw(tool_name, result)
        if not tool_result.ok:
            error = tool_result.error or "tool_error"
            logger.warning("⚠️ [{}] Tool execution returned a failure result: {}", self.name, error)
            _emit_node_runtime_event(
                self,
                "tool.call.failed",
                payload={
                    "tool": tool_name,
                    "mode": "workflow",
                    "error": error,
                },
            )
            trace_emit("tool.result", {
                "node": self.name,
                "tool": tool_name,
                "ok": False,
                "error": error,
            })
            self._emit_bt_skill_node_io(
                stage="completed",
                ok=False,
                tool_name=tool_name,
                input_value=tool_args,
                output_value=None,
                error=error,
            )
            return Status.FAILURE

        payload = tool_result.output
        self.write_port("output", payload, state_key=self.output_key)
        _emit_node_runtime_event(
            self,
            "tool.call.completed",
            payload={
                "tool": tool_name,
                "mode": "workflow",
                "result_data": self._event_result_payload(payload),
                "result_preview": self._trace_preview_payload(payload),
            },
        )

        trace_emit("tool.result", {
            "node": self.name,
            "tool": tool_name,
            "ok": True,
            "result": payload,
        })
        self._emit_bt_skill_node_io(
            stage="completed",
            ok=True,
            tool_name=tool_name,
            input_value=tool_args,
            output_value=payload,
            error=None,
        )
        return Status.SUCCESS
    except Exception as e:
        await _notify_tool_error_hooks(
            self,
            tool_name=getattr(self.tool, "name", type(self.tool).__name__),
            tool_args=locals().get("tool_args"),
            error=e,
            mode="workflow",
        )
        write_runtime_termination(
            self.state_manager,
            build_hook_error_termination(e, node_name=self.name),
        )
        logger.warning("⚠️ [{}] Tool execution failed: {}", self.name, e)
        _emit_node_runtime_event(
            self,
            "tool.call.failed",
            payload={
                "tool": getattr(self.tool, "name", type(self.tool).__name__),
                "mode": "workflow",
                "error": str(e),
            },
        )
        trace_emit("tool.result", {
            "node": self.name,
            "tool": getattr(self.tool, "name", type(self.tool).__name__),
            "ok": False,
            "error": str(e),
        })
        self._emit_bt_skill_node_io(
            stage="completed",
            ok=False,
            tool_name=getattr(self.tool, "name", type(self.tool).__name__),
            input_value=locals().get("tool_args"),
            output_value=None,
            error=str(e),
        )
        return Status.FAILURE

SkillNode

SkillNode(
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Any = None,
    skills_dir: str | Path | None = None,
    enabled_skills: list[str] | None = None,
    explicit_enabled_skills: bool = False,
    skill_files: list[str | Path] | None = None,
    messages: Sequence[Any] | None = None,
    tools: Optional[list[Tool]] = None,
    output_schema: Optional[dict[str, Any]] = None,
    config: SkillNodeConfig | None = None,
    tool_result_policy: ToolResultPolicy | None = None,
    tool_result_policy_scope: ToolResultPolicyScope
    | None = None,
)

Bases: FlattenedAgentNode

Skill-oriented agent node that loads and runs one or more skills.

SkillNode bridges file-based skill definitions into a runnable node. A skill may be prompt-driven or behavior-tree-driven, and the node handles skill discovery, prompt composition, tool exposure, and normalized output persistence.

属性:

名称 类型 描述
namespace

Optional state namespace used for this node's execution.

state_manager

Shared runtime state manager injected by the runner.

ctx

Runtime dependency context injected by the runner.

Initialize a skill-driven agent node.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
model_client Any

Explicit model client override for prompt-driven skills.

None
skills_dir str | Path | None

Directory containing SKILL.md files.

None
enabled_skills list[str] | None

Optional skill names to load from skills_dir.

None
explicit_enabled_skills bool

Whether an explicitly provided empty enabled_skills list should stay empty instead of loading all available catalog skills.

False
skill_files list[str | Path] | None

Explicit SKILL.md paths to load.

None
messages Sequence[Any] | None

Optional fixed messages injected into skill execution context.

None
tools Optional[list[Tool]]

Extra tools exposed alongside skill-declared tools.

None
output_schema Optional[dict[str, Any]]

Explicit output schema override.

None
config SkillNodeConfig | None

Skill node configuration and constraints.

None

引发:

类型 描述
ValueError

If skill_prompt_mode is not supported.

源代码位于: jianmu/node/builtin/skill.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Any = None,
    skills_dir: str | Path | None = None,
    enabled_skills: list[str] | None = None,
    explicit_enabled_skills: bool = False,
    skill_files: list[str | Path] | None = None,
    messages: Sequence[Any] | None = None,
    tools: Optional[list[Tool]] = None,
    output_schema: Optional[dict[str, Any]] = None,
    config: SkillNodeConfig | None = None,
    tool_result_policy: ToolResultPolicy | None = None,
    tool_result_policy_scope: ToolResultPolicyScope | None = None,
):
    """Initialize a skill-driven agent node.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        model_client: Explicit model client override for prompt-driven skills.
        skills_dir: Directory containing ``SKILL.md`` files.
        enabled_skills: Optional skill names to load from ``skills_dir``.
        explicit_enabled_skills: Whether an explicitly provided empty
            ``enabled_skills`` list should stay empty instead of loading all
            available catalog skills.
        skill_files: Explicit ``SKILL.md`` paths to load.
        messages: Optional fixed messages injected into skill execution context.
        tools: Extra tools exposed alongside skill-declared tools.
        output_schema: Explicit output schema override.
        config: Skill node configuration and constraints.

    Raises:
        ValueError: If ``skill_prompt_mode`` is not supported.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)

    resolved_config = config or SkillNodeConfig()
    self._skills_dir = Path(skills_dir) if skills_dir is not None else None
    self._enabled_skills = list(enabled_skills or [])
    self._explicit_enabled_skills = bool(explicit_enabled_skills)
    self._skill_files = [Path(p) for p in (skill_files or [])]

    self._seed_messages = normalize_messages(messages, fallback_role="user")
    self._explicit_model_client = (
        require_model_client(
            model_client,
            source=f"{self.__class__.__name__}.model_client",
        )
        if model_client is not None
        else None
    )
    self._tools = list(tools or [])
    self._push_to_chat = resolved_config.push_to_chat
    self._use_history = resolved_config.use_history
    self._history_limit = resolved_config.history_limit
    self._output_schema = output_schema
    self._tool_result_policy = tool_result_policy
    self._tool_result_policy_scope = tool_result_policy_scope
    self._keys = resolved_config.keys.model_copy(deep=True)
    base_react = resolved_config.react_config
    self._react_config = base_react.model_copy(deep=True) if base_react is not None else ReActConfig()
    self._react_config.keys = self._keys.model_copy(deep=True)
    self._constraints = deepcopy(resolved_config.constraints) if resolved_config.constraints is not None else Constraints()
    self._base_system_prompt = self._react_config.system_prompt
    if resolved_config.skill_prompt_mode not in {"summary", "full"}:
        raise ValueError("skill_prompt_mode must be one of: 'summary', 'full'")
    self._skill_prompt_mode: Literal["summary", "full"] = resolved_config.skill_prompt_mode

    self._loaded_skills: list[Skill] = []

update

update() -> Status

Run the compiled skill subtree and publish its formal outputs outward.

源代码位于: jianmu/node/builtin/skill.py
def update(self) -> Status:
    """Run the compiled skill subtree and publish its formal outputs outward."""
    status = super().update()
    if status != Status.RUNNING and self.state_manager is not None:
        direct_error = self.state_manager.get(
            _DIRECT_BT_ERROR_KEY,
            namespace=self.namespace,
            default=None,
        )
        final_text = self.state_manager.get(
            self._keys.final_answer,
            namespace=self.namespace,
            default=None,
        )
        if final_text:
            self.state_manager.update({self._keys.final_answer: final_text})
        skill_result = self.state_manager.get(
            self._keys.skill_result,
            namespace=self.namespace,
            default=None,
        )
        if skill_result is not None:
            self.state_manager.update({self._keys.skill_result: skill_result})
        if direct_error:
            self.state_manager.update(
                {_DIRECT_BT_ERROR_KEY: None},
                namespace=self.namespace,
                signal=False,
            )
            logger.error("❌ [{}] {}", self.name, direct_error)
            return Status.FAILURE
    return status

Log

Log(
    name: str | None = None,
    *,
    namespace: str | None = None,
    message: str = "",
)

Bases: Node

Log a message to console (and broadcast to Studio if configured).

属性:

名称 类型 描述
_broadcast_callback Optional[Callable[[str, str], None]]

Optional callback used to mirror logs externally.

message

Static log message emitted when the node ticks.

Create a node that logs one static message.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace.

None
message str

Static string to emit when the node ticks.

''
源代码位于: jianmu/node/builtin/utility.py
def __init__(self, name: str | None = None, *, namespace: str | None = None, message: str = ""):
    """Create a node that logs one static message.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace.
        message: Static string to emit when the node ticks.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.message = message

update

update() -> Status

Emit the configured log message and succeed.

返回:

类型 描述
Status

Always Status.SUCCESS.

源代码位于: jianmu/node/builtin/utility.py
def update(self) -> Status:
    """Emit the configured log message and succeed.

    Returns:
        Always ``Status.SUCCESS``.
    """
    log_msg = f"[{self.name}] {self.message}"
    logger.info("📝 [Log] {}: {}", self.name, self.message)
    if Log._broadcast_callback:
        Log._broadcast_callback("log", log_msg)
    return Status.SUCCESS

Wait

Wait(
    name: str | None = None,
    *,
    namespace: str | None = None,
    duration: float = 1.0,
)

Bases: AsyncNode

Wait for a specified duration, then return SUCCESS.

属性:

名称 类型 描述
duration

Number of seconds waited before succeeding.

Create an async wait node for a fixed duration.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace.

None
duration float

Number of seconds to wait before returning SUCCESS.

1.0
源代码位于: jianmu/node/builtin/utility.py
def __init__(self, name: str | None = None, *, namespace: str | None = None, duration: float = 1.0):
    """Create an async wait node for a fixed duration.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace.
        duration: Number of seconds to wait before returning SUCCESS.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.duration = float(duration)

update_async async

update_async() -> Status

Sleep for the configured duration and then succeed.

返回:

类型 描述
Status

Always Status.SUCCESS after the sleep completes.

源代码位于: jianmu/node/builtin/utility.py
async def update_async(self) -> Status:
    """Sleep for the configured duration and then succeed.

    Returns:
        Always ``Status.SUCCESS`` after the sleep completes.
    """
    logger.debug("⏳ [{}] Waiting {}s...", self.name, self.duration)
    await asyncio.sleep(self.duration)
    return Status.SUCCESS

Timeout

Timeout(
    name: str | None = None,
    *,
    namespace: str | None = None,
    child: Behaviour,
    duration: float,
)

Bases: Decorator

Async-friendly timeout decorator.

属性:

名称 类型 描述
duration

Timeout window in seconds.

expiry_time

Monotonic deadline for the current entry.

Wrap a child with an async-friendly timeout policy.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace.

None
child Behaviour

Child behaviour node to wrap.

必需
duration float

Timeout in seconds after which the child is cancelled.

必需
源代码位于: jianmu/node/builtin/utility.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    child: py_trees.behaviour.Behaviour,
    duration: float,
):
    """Wrap a child with an async-friendly timeout policy.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace.
        child: Child behaviour node to wrap.
        duration: Timeout in seconds after which the child is cancelled.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(child=child, name=resolved_name, namespace=namespace)
    self.duration = duration
    self.expiry_time = None
    self._restored_remaining_seconds: float | None = None

dump_resume_state

dump_resume_state() -> dict

Return remaining timeout budget for checkpoint restore.

源代码位于: jianmu/node/builtin/utility.py
def dump_resume_state(self) -> dict:
    """Return remaining timeout budget for checkpoint restore."""
    self._write_timeout_progress()
    remaining = self._remaining_seconds()
    if remaining is None:
        return {}
    return {"remaining_seconds": remaining}

restore_resume_state

restore_resume_state(data: dict) -> None

Restore pending timeout budget from checkpoint metadata.

源代码位于: jianmu/node/builtin/utility.py
def restore_resume_state(self, data: dict) -> None:
    """Restore pending timeout budget from checkpoint metadata."""
    try:
        remaining = float(dict(data or {}).get("remaining_seconds"))
    except Exception:
        self._restored_remaining_seconds = None
        return
    self._restored_remaining_seconds = max(0.0, remaining)

initialise

initialise() -> None

Start a fresh timeout window for the current entry.

源代码位于: jianmu/node/builtin/utility.py
def initialise(self) -> None:
    """Start a fresh timeout window for the current entry."""
    if self.consume_restored_entry() == "RUNNING":
        progress_remaining = self._read_timeout_progress()
        if progress_remaining is not None:
            duration = progress_remaining
        elif self._restored_remaining_seconds is not None:
            duration = self._restored_remaining_seconds
            self._restored_remaining_seconds = None
        else:
            duration = self.duration
    else:
        duration = self.duration
    self.expiry_time = time.monotonic() + duration
    self.feedback_message = ""

update

update() -> Status

Fail when the timeout expires; otherwise mirror child status.

返回:

类型 描述
Status

Status.FAILURE when the deadline passes, otherwise the current

Status

status of the child behaviour.

源代码位于: jianmu/node/builtin/utility.py
def update(self) -> Status:
    """Fail when the timeout expires; otherwise mirror child status.

    Returns:
        ``Status.FAILURE`` when the deadline passes, otherwise the current
        status of the child behaviour.
    """
    now = time.monotonic()
    if (
        self.decorated.status == Status.RUNNING
        and self.expiry_time
        and now > self.expiry_time
    ):
        self.feedback_message = f"\u23f0 Timed out after {self.duration}s"
        if isinstance(self.decorated, AsyncNode):
            if self.decorated.async_task and not self.decorated.async_task.done():
                self.decorated.async_task.cancel()
        self.decorated.stop(Status.INVALID)
        return Status.FAILURE

    if self.decorated.status == Status.RUNNING and self.expiry_time:
        remaining = max(0.0, self.expiry_time - now)
        self._write_timeout_progress()
        self.feedback_message = f"time still ticking ... [remaining: {remaining:.2f}s]"
    else:
        self.feedback_message = "child finished before timeout triggered"

    return self.decorated.status

StateCondition

StateCondition(
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    key: Optional[str] = None,
    op: str = "truthy",
    value: Any = None,
    predicate: Optional[Callable[[Any], bool]] = None,
)

Bases: Node

Generic state condition node.

Supports two evaluation modes: 1. Predicate mode: call predicate(state) and convert the result to bool. 2. Declarative mode: compare state[key] with op and value.

Returns SUCCESS when the condition passes, otherwise FAILURE.

属性:

名称 类型 描述
key

State field key used for declarative comparisons.

op

Comparison operator used for declarative mode.

value

Reference value used for declarative comparisons.

predicate

Optional predicate callable evaluated against the full state.

Configure a predicate- or comparison-based state condition.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace Optional[str]

Optional state namespace.

None
key Optional[str]

State field key to read for declarative comparisons.

None
op str

Comparison operator ("truthy", "falsy", "==", "!=", ">", ">=", "<", "<=").

'truthy'
value Any

Reference value for declarative comparisons.

None
predicate Optional[Callable[[Any], bool]]

Optional callable receiving the full state object; takes precedence over declarative comparisons when set.

None
源代码位于: jianmu/node/builtin/condition.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: Optional[str] = None,
    key: Optional[str] = None,
    op: str = "truthy",  # truthy, falsy, ==, !=, >, <, >=, <=
    value: Any = None,
    predicate: Optional[Callable[[Any], bool]] = None,
):
    """Configure a predicate- or comparison-based state condition.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace.
        key: State field key to read for declarative comparisons.
        op: Comparison operator (``"truthy"``, ``"falsy"``, ``"=="``,
            ``"!="``, ``">"``, ``">="``, ``"<"``, ``"<="``).
        value: Reference value for declarative comparisons.
        predicate: Optional callable receiving the full state object;
            takes precedence over declarative comparisons when set.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(resolved_name, namespace=namespace)
    self.key = key
    self.op = op
    self.value = value
    self.predicate = predicate

update

update() -> Status

Evaluate the configured condition against current state.

返回:

类型 描述
Status

Status.SUCCESS when the condition passes, Status.FAILURE

Status

when it fails or the state manager is not bound.

源代码位于: jianmu/node/builtin/condition.py
def update(self) -> Status:
    """Evaluate the configured condition against current state.

    Returns:
        ``Status.SUCCESS`` when the condition passes, ``Status.FAILURE``
        when it fails or the state manager is not bound.
    """
    if self.state_manager is None:
        return Status.FAILURE

    state = self.state_manager.get()

    # Predicate mode takes precedence over declarative comparisons.
    if self.predicate is not None:
        passed = bool(self.predicate(state))
        return Status.SUCCESS if passed else Status.FAILURE

    if self.key is None:
        return Status.FAILURE

    actual = self.read_port(self.key, default=None, state_key=self.key)
    passed = self._check(actual)
    return Status.SUCCESS if passed else Status.FAILURE

EvaluationNode

EvaluationNode(
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Optional[ModelClient] = None,
    input_keys: Optional[list[str]] = None,
    output_keys: Optional[
        EvaluationStateKeys | dict
    ] = None,
    append_to_messages: bool = True,
    feedback_role: str = "user",
    config: Optional[AgentLLMConfig] = None,
)

Bases: AsyncNode

Evaluator node for actor-evaluator workflows.

Reads configured input fields from state, asks the model for a structured score/reflection result, and writes outputs back through EvaluationStateKeys.

属性:

名称 类型 描述
_explicit_model_client

Explicit evaluator model-client override.

config

Agent LLM configuration used for the evaluator call.

output_keys

State-key mapping used for score/reflection outputs.

input_keys

State field names included in the evaluator prompt.

append_to_messages

Whether evaluator feedback is appended to history.

feedback_role

Role used when appending feedback messages.

Configure evaluator inputs, outputs, and provider settings.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
model_client Optional[ModelClient]

Explicit model client override for the evaluator call.

None
input_keys Optional[list[str]]

State field names to include in the evaluator prompt.

None
output_keys Optional[EvaluationStateKeys | dict]

State keys that receive the score, reflection, and best-answer outputs.

None
append_to_messages bool

Whether to append the feedback message to the conversation history.

True
feedback_role str

Role used when appending the feedback message.

'user'
config Optional[AgentLLMConfig]

Agent LLM configuration for the evaluator model call.

None
源代码位于: jianmu/node/builtin/evaluation.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Optional[ModelClient] = None,
    input_keys: Optional[list[str]] = None,
    output_keys: Optional[EvaluationStateKeys | dict] = None,
    append_to_messages: bool = True,
    feedback_role: str = "user",
    config: Optional[AgentLLMConfig] = None,
):
    """Configure evaluator inputs, outputs, and provider settings.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        model_client: Explicit model client override for the evaluator call.
        input_keys: State field names to include in the evaluator prompt.
        output_keys: State keys that receive the score, reflection, and
            best-answer outputs.
        append_to_messages: Whether to append the feedback message to the
            conversation history.
        feedback_role: Role used when appending the feedback message.
        config: Agent LLM configuration for the evaluator model call.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(resolved_name, namespace=namespace)
    self._explicit_model_client = model_client
    if config is None:
        config = AgentLLMConfig(model=get_config().models.evaluate)
    self.config = config.model_copy(deep=True)
    resolved_keys = output_keys
    if isinstance(resolved_keys, dict):
        resolved_keys = EvaluationStateKeys(**resolved_keys)
    self.output_keys = (
        resolved_keys.model_copy(deep=True) if resolved_keys is not None else EvaluationStateKeys()
    )
    default_input_keys = [self.config.keys.messages, self.output_keys.best_answer_source]
    self.input_keys = list(dict.fromkeys(input_keys or default_input_keys))
    self.append_to_messages = append_to_messages
    self.feedback_role = feedback_role

update_async async

update_async() -> Status

Call the evaluator model and persist score/reflection outputs.

返回:

类型 描述
Status

Status.SUCCESS when a valid score and reflection are extracted,

Status

Status.FAILURE when the provider is missing, state is not

Status

available, or the model response cannot be parsed.

源代码位于: jianmu/node/builtin/evaluation.py
async def update_async(self) -> Status:
    """Call the evaluator model and persist score/reflection outputs.

    Returns:
        ``Status.SUCCESS`` when a valid score and reflection are extracted,
        ``Status.FAILURE`` when the provider is missing, state is not
        available, or the model response cannot be parsed.
    """
    if self.state_manager is None:
        logger.error("[{}] StateManager not injected.", self.name)
        return Status.FAILURE

    state = self.state_manager.get()
    content = self._build_prompt(state)

    tools_schema = [{
        "type": "function",
        "function": {
            "name": "evaluate_result",
            "description": "Output the evaluation score and reflection.",
            "parameters": EvaluationResult.model_json_schema(),
        },
    }]

    messages = self._resolve_context_builder().build(
        local_state={"messages": [Message(role="user", content=content)]},
        global_state=state,
        ctx=self.ctx,
    )

    model_config = ModelConfig(
        model=self.config.model,
        temperature=self.config.temperature,
        max_tokens=self.config.max_tokens,
        top_p=self.config.top_p,
        top_k=self.config.top_k or 40,
        timeout=self.config.timeout,
        max_retries=self.config.resilience.max_retries,
        fallback_model=self.config.resilience.fallback_model,
        disable_fallback=not self.config.resilience.enable_fallback,
        **self.config.extra_params,
        tools=tools_schema,
        tool_choice={"type": "function", "function": {"name": "evaluate_result"}},
    )

    try:
        model_client = self.model_client or ModelClient.resolve()
        response, _ = await model_client.invoke(
            messages=messages,
            config=model_config,
            stream=False,
            trace_node=self.name,
        )
        if response.tool_calls and len(response.tool_calls) > 0:
            extracted = extract_tool_call_from_dict(response.tool_calls[0])
            if not extracted:
                logger.error("[{}] Failed to extract structured output from LLM response.", self.name)
                return Status.FAILURE

            args = extracted.arguments
            if isinstance(args, str):
                args = json.loads(args)

            score = float(args.get("score", 0.0))
            reflection = args.get("reflection", "")

            updates = {}
            if self.output_keys.score:
                updates["score"] = score
            if self.output_keys.reflection:
                updates["reflection"] = reflection

            highest_score_key = self.output_keys.highest_score
            best_answer_key = self.output_keys.best_answer
            best_answer_source_key = self.output_keys.best_answer_source
            if highest_score_key and best_answer_key and best_answer_source_key:
                highest_so_far = self._read_with_fallback(state, highest_score_key, default=0.0)
                best_answer_source = self._read_with_fallback(state, best_answer_source_key, default=None)
                if score > highest_so_far and best_answer_source is not None:
                    updates["highest_score"] = score
                    updates["best_answer"] = best_answer_source

            if self.append_to_messages:
                feedback = f"Score: {score}/10\nReflection: {reflection}"
                self.append_port_messages(
                    "messages",
                    [
                    Message(role=self.feedback_role, content=feedback)
                    ],
                    state_key=self.config.keys.messages,
                    signal=False,
                )

            self.write_ports(
                updates,
                state_keys={
                    "score": self.output_keys.score,
                    "reflection": self.output_keys.reflection,
                    "highest_score": highest_score_key,
                    "best_answer": best_answer_key,
                },
            )

            trace_emit(
                "evaluator.result",
                {
                    "node": self.name,
                    "score": score,
                    "reflection": reflection,
                },
            )
            logger.debug("[{}] Generated Evaluation - Score: {}, Reflection: {}", self.name, score, reflection)
            return Status.SUCCESS

        logger.error("[{}] No tool calls returned from LLM for evaluation.", self.name)
        return Status.FAILURE

    except Exception as exc:
        logger.error("[{}] Evaluation failed: {}", self.name, exc)
        return Status.FAILURE

SpawnAgent

SpawnAgent(
    name: str | None = None,
    *,
    namespace: str | None = None,
    runtime: Optional[Any] = None,
    role: Optional[str] = None,
    task: Optional[str] = None,
    parent_id: Optional[str] = None,
    constraints: Optional[object] = None,
    role_key: str = "role",
    task_key: str = "task",
    parent_id_key: str = "agent_id",
    constraints_key: str = "constraints",
    mode: str = "detached",
    output_key: str = "spawned_agent_id",
    reuse_if_present: bool = False,
)

Bases: _SwarmRuntimeNode

Behavior tree node to spawn a new agent in the swarm runtime.

属性:

名称 类型 描述
role

Static role name override.

role_key

State key used to resolve the role dynamically.

task

Static task text override.

task_key

State key used to resolve the task dynamically.

mode

Spawn mode passed to the runtime.

constraints

Static execution constraints override.

constraints_key

State key used to resolve constraints dynamically.

runtime

Explicit swarm runtime reference.

parent_id

Static parent agent id override.

parent_id_key

State key used to resolve the parent agent id.

output_key

State key where the spawned agent id is written.

reuse_if_present

Whether an existing output agent id should be reused.

Initialize a node that spawns swarm agents.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
runtime Optional[Any]

Explicit swarm runtime reference; falls back to state if omitted.

None
role Optional[str]

Static role name to spawn; overrides role_key lookup.

None
task Optional[str]

Static task string; overrides task_key lookup.

None
parent_id Optional[str]

Static parent agent id; overrides parent_id_key lookup.

None
constraints Optional[object]

Static execution constraints; overrides constraints_key lookup.

None
role_key str

State key used to read the role when role is not set.

'role'
task_key str

State key used to read the task when task is not set.

'task'
parent_id_key str

State key for the parent agent id.

'agent_id'
constraints_key str

State key for runtime constraints.

'constraints'
mode str

Agent spawn mode (e.g. "detached").

'detached'
output_key str

State key where the new agent id is written.

'spawned_agent_id'
reuse_if_present bool

When True, reuse a previously written agent id from output_key instead of spawning a new agent.

False
源代码位于: jianmu/node/builtin/swarm.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    runtime: Optional[Any] = None,
    role: Optional[str] = None,
    task: Optional[str] = None,
    parent_id: Optional[str] = None,
    constraints: Optional[object] = None,
    role_key: str = "role",
    task_key: str = "task",
    parent_id_key: str = "agent_id",
    constraints_key: str = "constraints",
    mode: str = "detached",
    output_key: str = "spawned_agent_id",
    reuse_if_present: bool = False,
):
    """Initialize a node that spawns swarm agents.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        runtime: Explicit swarm runtime reference; falls back to state if
            omitted.
        role: Static role name to spawn; overrides ``role_key`` lookup.
        task: Static task string; overrides ``task_key`` lookup.
        parent_id: Static parent agent id; overrides ``parent_id_key``
            lookup.
        constraints: Static execution constraints; overrides
            ``constraints_key`` lookup.
        role_key: State key used to read the role when ``role`` is not set.
        task_key: State key used to read the task when ``task`` is not set.
        parent_id_key: State key for the parent agent id.
        constraints_key: State key for runtime constraints.
        mode: Agent spawn mode (e.g. ``"detached"``).
        output_key: State key where the new agent id is written.
        reuse_if_present: When True, reuse a previously written agent id
            from ``output_key`` instead of spawning a new agent.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.role = role
    self.role_key = role_key
    self.task = task
    self.task_key = task_key
    self.mode = mode
    self.constraints = constraints
    self.constraints_key = constraints_key
    self.runtime = runtime
    self.parent_id = parent_id
    self.parent_id_key = parent_id_key
    self.output_key = output_key
    self.reuse_if_present = bool(reuse_if_present)

update_async async

update_async() -> Status

Spawn an agent and write its identifier to state.

When reuse_if_present is enabled, the node first checks the resolved output_key state slot. If a non-empty agent id is already present, that id is preserved and no new agent is spawned.

返回:

类型 描述
Status

Status.SUCCESS when the agent id is resolved or spawned,

Status

Status.FAILURE when the runtime or role cannot be resolved.

源代码位于: jianmu/node/builtin/swarm.py
async def update_async(self) -> Status:
    """Spawn an agent and write its identifier to state.

    When ``reuse_if_present`` is enabled, the node first checks the
    resolved ``output_key`` state slot. If a non-empty agent id is already
    present, that id is preserved and no new agent is spawned.

    Returns:
        ``Status.SUCCESS`` when the agent id is resolved or spawned,
        ``Status.FAILURE`` when the runtime or role cannot be resolved.
    """
    runtime = self._resolve_runtime(required_methods=("spawn",))
    if runtime is None:
        return Status.FAILURE

    if self.reuse_if_present and self.output_key:
        resolved_output_key = self.resolve_output_port("spawned_agent_id", state_key=self.output_key)
        existing_agent_id = self.read_state(resolved_output_key, default=None)
        if isinstance(existing_agent_id, str) and existing_agent_id.strip():
            self.write_port("spawned_agent_id", existing_agent_id, state_key=self.output_key)
            return Status.SUCCESS

    role = self.role or self._read_state_value(self.role_key)
    if not role:
        return Status.FAILURE

    task = self.task
    if task is None:
        task = self._read_state_value(self.task_key)

    parent_id = self.parent_id
    if parent_id is None:
        parent_id = self._read_state_value(self.parent_id_key)

    constraints = self.constraints
    if constraints is None:
        constraints = self._read_state_value(self.constraints_key)

    recovery_task = self.recovery.task(
        kind="swarm.spawn",
        input={
            "role": role,
            "task": task,
            "parent_id": parent_id,
            "mode": self.mode,
            "constraints": self._record_value(constraints),
            "output_key": self.output_key,
        },
    )
    completed = self._completed_task(recovery_task)
    if completed is not None:
        result = completed.get("result") if isinstance(completed.get("result"), dict) else {}
        agent_id = str(result.get("agent_id") or completed.get("agent_id") or "")
        if not agent_id:
            return Status.FAILURE
        if self.output_key:
            self.write_port("spawned_agent_id", agent_id, state_key=self.output_key)
        return Status.SUCCESS
    if self._blocked_uncertain_task(recovery_task):
        return Status.FAILURE

    self.recovery.tasks.begin(recovery_task)
    try:
        agent_id = runtime.spawn(
            role=role,
            task=task,
            parent_id=parent_id,
            mode=self.mode,
            constraints=constraints,
        )
    except Exception as exc:
        self.recovery.tasks.mark_failed(
            str(recovery_task.get("task_id") or ""),
            error={"message": str(exc), "type": exc.__class__.__name__},
        )
        raise
    if self.output_key:
        self.write_port("spawned_agent_id", agent_id, state_key=self.output_key)
    self.recovery.tasks.mark_completed(
        str(recovery_task.get("task_id") or ""),
        result={"agent_id": agent_id},
        extra={
            "kind": "swarm.spawn",
            "node_locator": self.node_locator,
            "owner": recovery_task.get("owner"),
            "input": recovery_task.get("input"),
            "input_hash": recovery_task.get("input_hash"),
            "agent_id": agent_id,
        },
    )
    return Status.SUCCESS

SendMessage

SendMessage(
    name: str | None = None,
    *,
    namespace: str | None = None,
    runtime: Optional[Any] = None,
    to_agent_id: Optional[str] = None,
    group_id: Optional[str] = None,
    topic: Optional[str] = None,
    content: Optional[str] = None,
    sender_id: Optional[str] = None,
    to_agent_key: str = "to_agent_id",
    group_key: str = "group_id",
    topic_key: str = "topic",
    content_key: str = "message",
    sender_key: str = "agent_id",
)

Bases: _SwarmRuntimeNode

Behavior tree node to send a message to another agent or group.

属性:

名称 类型 描述
to_agent_id

Static direct-message recipient agent id.

to_agent_key

State key used to resolve the recipient dynamically.

group_id

Static group channel id.

group_key

State key used to resolve the group dynamically.

topic

Static topic channel name.

topic_key

State key used to resolve the topic dynamically.

content

Static message body override.

content_key

State key used to resolve the message body dynamically.

sender_id

Static sender id override.

sender_key

State key used to resolve the sender dynamically.

runtime

Explicit swarm runtime reference.

Initialize a node that sends swarm messages.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
runtime Optional[Any]

Explicit swarm runtime reference; falls back to state.

None
to_agent_id Optional[str]

Static direct-message recipient agent id.

None
group_id Optional[str]

Static group channel id.

None
topic Optional[str]

Static topic channel name.

None
content Optional[str]

Static message body; overrides content_key lookup.

None
sender_id Optional[str]

Static sender id; overrides sender_key lookup.

None
to_agent_key str

State key for the recipient agent id.

'to_agent_id'
group_key str

State key for the group id.

'group_id'
topic_key str

State key for the topic name.

'topic'
content_key str

State key for the message body.

'message'
sender_key str

State key for the sender id.

'agent_id'
源代码位于: jianmu/node/builtin/swarm.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    runtime: Optional[Any] = None,
    to_agent_id: Optional[str] = None,
    group_id: Optional[str] = None,
    topic: Optional[str] = None,
    content: Optional[str] = None,
    sender_id: Optional[str] = None,
    to_agent_key: str = "to_agent_id",
    group_key: str = "group_id",
    topic_key: str = "topic",
    content_key: str = "message",
    sender_key: str = "agent_id",
):
    """Initialize a node that sends swarm messages.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        runtime: Explicit swarm runtime reference; falls back to state.
        to_agent_id: Static direct-message recipient agent id.
        group_id: Static group channel id.
        topic: Static topic channel name.
        content: Static message body; overrides ``content_key`` lookup.
        sender_id: Static sender id; overrides ``sender_key`` lookup.
        to_agent_key: State key for the recipient agent id.
        group_key: State key for the group id.
        topic_key: State key for the topic name.
        content_key: State key for the message body.
        sender_key: State key for the sender id.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)
    self.to_agent_id = to_agent_id
    self.to_agent_key = to_agent_key
    self.group_id = group_id
    self.group_key = group_key
    self.topic = topic
    self.topic_key = topic_key
    self.content = content
    self.content_key = content_key
    self.sender_id = sender_id
    self.sender_key = sender_key
    self.runtime = runtime

update_async async

update_async() -> Status

Send a routed swarm message through the runtime.

返回:

类型 描述
Status

Status.SUCCESS when the message is delivered, Status.FAILURE

Status

when the runtime, content, or routing target cannot be resolved.

源代码位于: jianmu/node/builtin/swarm.py
async def update_async(self) -> Status:
    """Send a routed swarm message through the runtime.

    Returns:
        ``Status.SUCCESS`` when the message is delivered, ``Status.FAILURE``
        when the runtime, content, or routing target cannot be resolved.
    """
    runtime = self._resolve_runtime(required_methods=("send_message", "emit"))
    if runtime is None:
        return Status.FAILURE

    content = self.content
    if content is None:
        content = self._read_state_value(self.content_key)
    if content is None:
        return Status.FAILURE

    to_agent_id = self.to_agent_id or self._read_state_value(self.to_agent_key)
    group_id = self.group_id or self._read_state_value(self.group_key)
    topic = self.topic or self._read_state_value(self.topic_key)
    if not to_agent_id and not group_id and not topic:
        return Status.FAILURE

    sender_id = self.sender_id or self._read_state_value(self.sender_key)
    delivery_mode = "send_message" if hasattr(runtime, "send_message") else "emit"
    recovery_task = self.recovery.task(
        kind="swarm.send_message",
        input={
            "sender_id": sender_id,
            "content": str(content),
            "to_agent_id": to_agent_id,
            "group_id": group_id,
            "topic": topic,
            "delivery_mode": delivery_mode,
        },
    )
    if self._completed_task(recovery_task) is not None:
        return Status.SUCCESS
    if self._blocked_uncertain_task(recovery_task):
        return Status.FAILURE

    self.recovery.tasks.begin(recovery_task)
    message_id = ""
    if hasattr(runtime, "send_message"):
        try:
            await runtime.send_message(
                sender_id=sender_id,
                content=str(content),
                to_agent_id=to_agent_id,
                group_id=group_id,
                topic=topic,
            )
        except Exception as exc:
            self.recovery.tasks.mark_failed(
                str(recovery_task.get("task_id") or ""),
                error={"message": str(exc), "type": exc.__class__.__name__},
            )
            raise
    else:
        message = Message(
            id=str(uuid.uuid4()),
            role="user",
            content=str(content),
            metadata={
                "sender_id": sender_id,
                "content_type": "text",
                "routing": {
                    "to_agent_id": to_agent_id,
                    "group_id": group_id,
                    "topic": topic,
                },
                "created_at": time.time(),
            },
        )
        from jianmu.swarm.types import AgentEvent
        try:
            runtime.emit(
                AgentEvent(
                    event_type="message_sent",
                    agent_id=sender_id,
                    payload={
                        "id": message.id,
                        "to": to_agent_id,
                        "group_id": group_id,
                        "topic": topic,
                        "content": message.content,
                        "content_type": "text",
                        "created_at": message.metadata.get("created_at"),
                        "metadata": message.metadata,
                    },
                )
            )
        except Exception as exc:
            self.recovery.tasks.mark_failed(
                str(recovery_task.get("task_id") or ""),
                error={"message": str(exc), "type": exc.__class__.__name__},
            )
            raise
        message_id = message.id
    self.recovery.tasks.mark_completed(
        str(recovery_task.get("task_id") or ""),
        result={"message_id": message_id},
        extra={
            "kind": "swarm.send_message",
            "node_locator": self.node_locator,
            "owner": recovery_task.get("owner"),
            "input": recovery_task.get("input"),
            "input_hash": recovery_task.get("input_hash"),
            "message_id": message_id,
        },
    )
    return Status.SUCCESS

WaitMessage

WaitMessage(
    name: str | None = None, *, namespace: str | None = None
)

Bases: Node

Block execution until inbound swarm messages are available.

Initialize a node that waits for inbound messages.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
源代码位于: jianmu/node/builtin/swarm.py
def __init__(self, name: str | None = None, *, namespace: str | None = None):
    """Initialize a node that waits for inbound messages.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)

update

update() -> Status

Succeed when incoming messages are available.

返回:

类型 描述
Status

Status.SUCCESS when ctx.incoming_messages is non-empty,

Status

otherwise Status.RUNNING.

源代码位于: jianmu/node/builtin/swarm.py
def update(self) -> Status:
    """Succeed when incoming messages are available.

    Returns:
        ``Status.SUCCESS`` when ``ctx.incoming_messages`` is non-empty,
        otherwise ``Status.RUNNING``.
    """
    incoming = getattr(self.ctx, "incoming_messages", None)
    if isinstance(incoming, list) and incoming:
        return Status.SUCCESS
    return Status.RUNNING

WaitForever

WaitForever(
    name: str | None = None, *, namespace: str | None = None
)

Bases: Node

Keep execution running indefinitely until interrupted from outside.

Initialize a node that never completes.

参数:

名称 类型 描述 默认
name str | None

Node name shown in traces and tree views.

None
namespace str | None

Optional state namespace for port resolution.

None
源代码位于: jianmu/node/builtin/swarm.py
def __init__(self, name: str | None = None, *, namespace: str | None = None):
    """Initialize a node that never completes.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)

update

update() -> Status

Keep the node running indefinitely.

返回:

类型 描述
Status

Always Status.RUNNING.

源代码位于: jianmu/node/builtin/swarm.py
def update(self) -> Status:
    """Keep the node running indefinitely.

    Returns:
        Always ``Status.RUNNING``.
    """
    return Status.RUNNING

SwarmNode

SwarmNode(
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Any = None,
    skills_dir: str | Path | None = None,
    enabled_skills: list[str] | None = None,
    explicit_enabled_skills: bool = False,
    skill_files: list[str | Path] | None = None,
    messages: Sequence[Any] | None = None,
    tools: Optional[list[Tool]] = None,
    output_schema: Optional[dict[str, Any]] = None,
    config: SkillNodeConfig | None = None,
    tool_result_policy: ToolResultPolicy | None = None,
    tool_result_policy_scope: ToolResultPolicyScope
    | None = None,
)

Bases: SkillNode

Behavior tree node that hosts swarm agent execution using skills.

源代码位于: jianmu/node/builtin/skill.py
def __init__(
    self,
    name: str | None = None,
    *,
    namespace: str | None = None,
    model_client: Any = None,
    skills_dir: str | Path | None = None,
    enabled_skills: list[str] | None = None,
    explicit_enabled_skills: bool = False,
    skill_files: list[str | Path] | None = None,
    messages: Sequence[Any] | None = None,
    tools: Optional[list[Tool]] = None,
    output_schema: Optional[dict[str, Any]] = None,
    config: SkillNodeConfig | None = None,
    tool_result_policy: ToolResultPolicy | None = None,
    tool_result_policy_scope: ToolResultPolicyScope | None = None,
):
    """Initialize a skill-driven agent node.

    Args:
        name: Node name shown in traces and tree views.
        namespace: Optional state namespace for port resolution.
        model_client: Explicit model client override for prompt-driven skills.
        skills_dir: Directory containing ``SKILL.md`` files.
        enabled_skills: Optional skill names to load from ``skills_dir``.
        explicit_enabled_skills: Whether an explicitly provided empty
            ``enabled_skills`` list should stay empty instead of loading all
            available catalog skills.
        skill_files: Explicit ``SKILL.md`` paths to load.
        messages: Optional fixed messages injected into skill execution context.
        tools: Extra tools exposed alongside skill-declared tools.
        output_schema: Explicit output schema override.
        config: Skill node configuration and constraints.

    Raises:
        ValueError: If ``skill_prompt_mode`` is not supported.
    """
    resolved_name = name or self.__class__.__name__
    super().__init__(name=resolved_name, namespace=namespace)

    resolved_config = config or SkillNodeConfig()
    self._skills_dir = Path(skills_dir) if skills_dir is not None else None
    self._enabled_skills = list(enabled_skills or [])
    self._explicit_enabled_skills = bool(explicit_enabled_skills)
    self._skill_files = [Path(p) for p in (skill_files or [])]

    self._seed_messages = normalize_messages(messages, fallback_role="user")
    self._explicit_model_client = (
        require_model_client(
            model_client,
            source=f"{self.__class__.__name__}.model_client",
        )
        if model_client is not None
        else None
    )
    self._tools = list(tools or [])
    self._push_to_chat = resolved_config.push_to_chat
    self._use_history = resolved_config.use_history
    self._history_limit = resolved_config.history_limit
    self._output_schema = output_schema
    self._tool_result_policy = tool_result_policy
    self._tool_result_policy_scope = tool_result_policy_scope
    self._keys = resolved_config.keys.model_copy(deep=True)
    base_react = resolved_config.react_config
    self._react_config = base_react.model_copy(deep=True) if base_react is not None else ReActConfig()
    self._react_config.keys = self._keys.model_copy(deep=True)
    self._constraints = deepcopy(resolved_config.constraints) if resolved_config.constraints is not None else Constraints()
    self._base_system_prompt = self._react_config.system_prompt
    if resolved_config.skill_prompt_mode not in {"summary", "full"}:
        raise ValueError("skill_prompt_mode must be one of: 'summary', 'full'")
    self._skill_prompt_mode: Literal["summary", "full"] = resolved_config.skill_prompt_mode

    self._loaded_skills: list[Skill] = []

update

update() -> Status

Run the skill node and persist a swarm success flag.

返回:

类型 描述
Status

The execution Status of the node.

源代码位于: jianmu/node/builtin/swarm.py
def update(self) -> Status:
    """Run the skill node and persist a swarm success flag.

    Returns:
        The execution Status of the node.
    """
    status = super().update()
    if status != Status.RUNNING and self.state_manager:
        self.state_manager.update({"swarm_success": status == Status.SUCCESS})
    return status

StateKeys

Bases: BaseModel

Single source of truth for state field naming conventions.

属性:

名称 类型 描述
messages str

State key for conversation history.

rounds str

State key for loop round counters.

actions str

State key for model-emitted actions or tool calls.

done str

State key indicating workflow completion.

text_output str

State key for latest plain-text node output.

final_answer str

State key for final user-facing answer text.

skill_result str

State key for structured skill outputs.

usage str

State key for aggregated model-usage metadata.

streaming_output str

State key for incremental streamed text.

tool_effects str

State key for aggregated tool-side-effect counters.

EvaluationStateKeys

Bases: BaseModel

State field routing for EvaluationNode outputs.

属性:

名称 类型 描述
score str | None

State key receiving the evaluator score.

reflection str | None

State key receiving the evaluator reflection text.

highest_score str | None

State key tracking the best score so far.

best_answer str | None

State key storing the best answer content.

best_answer_source str

State key pointing to the answer candidate being scored.

ReActConfig

Bases: BaseModel

ReAct (Reason + Act) loop configuration.

system_prompt configures the persona/base prompt. The ReAct protocol prompt and tool descriptions are added by the context builder. max_iterations is the main agent-loop round limit, not a low-level tree tick limit.

属性:

名称 类型 描述
model str

Primary model identifier used for ReAct turns.

temperature float

Sampling temperature for generation.

max_tokens Optional[int]

Optional output-token cap.

top_p Optional[float]

Optional nucleus-sampling parameter.

top_k Optional[int]

Optional top-k sampling parameter.

timeout float

Provider timeout in seconds.

stream bool

Whether streaming responses are requested.

max_budget_tokens Optional[int]

Optional budget cap shared with runtime accounting.

max_iterations int

Main ReAct loop round limit.

soft_landing bool

Whether soft-landing summary behavior is enabled.

resilience ModelResilienceConfig

Retry and fallback policy for provider calls.

system_prompt Optional[str]

Optional persona/base prompt override.

tool_choice Optional[Any]

Provider-specific tool selection directive applied only to the first model round of each ReAct run.

extra_params Dict[str, Any]

Additional provider-specific generation parameters.

keys StateKeys

State-key mapping used by the preset.

AgentLLMConfig

Bases: BaseModel

LLM generation parameters plus state routing keys.

system_prompt configures the persona/base prompt. Context presets may still append protocol, tools, skills, history, and other runtime blocks.

属性:

名称 类型 描述
model str

Primary model identifier used for inference.

temperature float

Sampling temperature for generation.

max_tokens Optional[int]

Optional output-token cap.

top_p Optional[float]

Optional nucleus-sampling parameter.

top_k Optional[int]

Optional top-k sampling parameter.

timeout float

Provider timeout in seconds.

stream bool

Whether streaming responses are requested.

max_budget_tokens Optional[int]

Optional budget cap shared with runtime accounting.

resilience ModelResilienceConfig

Retry and fallback policy for provider calls.

extra_params Dict[str, Any]

Additional provider-specific generation parameters.

system_prompt Optional[str]

Optional persona/base prompt override.

keys StateKeys

State-key mapping used by the node.

ToolExecutorConfig

Bases: BaseModel

Tool executor configuration.

属性:

名称 类型 描述
max_retries int

Maximum retry count for tool execution failures.

retry_backoff float

Backoff delay between retries.

timeout_s float | None

Optional timeout in seconds for each tool execution.

observation_format str

Format used for tool observations.

execution_mode Literal['auto', 'serial', 'parallel']

Tool scheduling mode: auto, serial, or parallel.

keys StateKeys

State-key mapping used by the executor.

PlanExecuteConfig

Bases: BaseModel

Plan-Execute(-Review) loop configuration.

属性:

名称 类型 描述
model str

Primary model identifier used across the preset.

temperature float

Sampling temperature for generation.

max_tokens Optional[int]

Optional output-token cap.

top_p Optional[float]

Optional nucleus-sampling parameter.

top_k Optional[int]

Optional top-k sampling parameter.

timeout float

Provider timeout in seconds.

max_budget_tokens Optional[int]

Optional budget cap shared with runtime accounting.

stream bool

Whether streaming responses are requested.

max_rounds int

Maximum plan/execute rounds before termination.

enable_review bool

Whether review steps are enabled.

review_threshold float

Score threshold for successful review.

plan_prompt Optional[str]

Optional custom planning prompt override.

execute_prompt Optional[str]

Optional custom execution prompt override.

review_prompt Optional[str]

Optional custom review prompt override.

plan_key str

State key used to store the plan.

score_key str

State key used to store review scores.

keys StateKeys

State-key mapping used by the preset.

node

node(
    _func: Optional[Callable[..., Any]] = None,
    *,
    name: Optional[str] = None,
    description: Optional[str] = None,
) -> (
    type[FunctionNode]
    | Callable[[Callable[..., Any]], type[FunctionNode]]
)

Decorator to wrap a function into a behaviour tree node class.

Can be used with or without arguments::

@node
def my_func(state): ...

@node(name="custom_name")
def my_func(state): ...

参数:

名称 类型 描述 默认
_func Optional[Callable[..., Any]]

The function to wrap when used without parentheses.

None
name Optional[str]

Optional node name override (defaults to the function name).

None
description Optional[str]

Optional node description (defaults to the function docstring).

None

返回:

类型 描述
type[FunctionNode] | Callable[[Callable[..., Any]], type[FunctionNode]]

A FunctionNode subclass that can be instantiated as a behaviour tree node.

源代码位于: jianmu/node/decorator.py
def node(
    _func: Optional[Callable[..., Any]] = None,
    *,
    name: Optional[str] = None,
    description: Optional[str] = None,
) -> type[FunctionNode] | Callable[[Callable[..., Any]], type[FunctionNode]]:
    """Decorator to wrap a function into a behaviour tree node class.

    Can be used with or without arguments::

        @node
        def my_func(state): ...

        @node(name="custom_name")
        def my_func(state): ...

    Args:
        _func: The function to wrap when used without parentheses.
        name: Optional node name override (defaults to the function name).
        description: Optional node description (defaults to the function docstring).

    Returns:
        A ``FunctionNode`` subclass that can be instantiated as a behaviour tree node.
    """
    def decorator(func: Callable[..., Any]) -> type[FunctionNode]:
        """Wrap a function as an instantiable behaviour node class.

        Args:
            func: The wrapped callable matching the node signature.
        """
        node_name, node_desc = _get_metadata(func, name, description)

        # We return a class that can be instantiated with (name, state)
        # to match the existing usage pattern in tests.
        class WrappedNode(FunctionNode):
            """Dynamic behavior tree node class wrapping a python function."""

            def __init__(
                self,
                inst_name: Optional[str] = None,
                state_manager: Optional[StateManager] = None,
            ):
                """Instantiate the generated node wrapper."""
                # If name is provided during instantiation, it overrides the decorator name
                super().__init__(inst_name or node_name, state_manager, func)
                self.description = node_desc

        WrappedNode.__name__ = f"Node_{func.__name__}"
        WrappedNode.__doc__ = node_desc
        return WrappedNode

    if _func is None:
        return decorator
    return decorator(_func)

create_react_node

create_react_node(
    name: str = "ReActAgent",
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    tools: Optional[Sequence[Tool]] = None,
    tool_providers: Optional[Sequence[ToolProvider]] = None,
    tool_provider_context: Optional[Any] = None,
    tool_runner: Optional[Any] = None,
    constraints: Optional[Any] = None,
    context_builder: Optional[
        ContextBuilderProtocol
    ] = None,
    config: Optional[ReActConfig] = None,
    tool_result_policy: Optional[ToolResultPolicy] = None,
    tool_result_policy_scope: Optional[
        ToolResultPolicyScope
    ] = None,
) -> Selector

Build a ReAct subtree with agent, tool, and completion nodes.

This preset assembles the common ReAct loop:

  1. ask an agent LLM node for the next step
  2. execute emitted tool calls
  3. stop once done becomes truthy in state

It is the quickest way to build a tool-using agent::

root = create_react_node(
    model_client=my_model_client,
    tools=[CalculatorTool()],
)

参数:

名称 类型 描述 默认
name str

Root node name for the preset subtree.

'ReActAgent'
namespace Optional[str]

Optional state namespace for isolation.

None
model_client ModelClient

Model client used by the agent node.

必需
tools Optional[Sequence[Tool]]

Static tools available to the agent.

None
tool_providers Optional[Sequence[ToolProvider]]

Optional providers that contribute dynamic tools at assembly time.

None
tool_provider_context Optional[Any]

Optional context passed to tool providers during collection.

None
tool_runner Optional[Any]

Optional explicit tool runner override.

None
constraints Optional[Any]

Optional execution/guard constraints shared with the tool executor.

None
context_builder Optional[ContextBuilderProtocol]

Optional prompt context builder passed to the agent LLM node.

None
config Optional[ReActConfig]

ReAct configuration. Defaults are used when omitted.

None
tool_result_policy Optional[ToolResultPolicy]

Optional result policy evaluated after tool execution and before observation commit.

None
tool_result_policy_scope Optional[ToolResultPolicyScope]

Optional predicate limiting the calls to which tool_result_policy applies. Without one, policy-enabled batches are conservatively serialized.

None

返回:

类型 描述
Selector

Configured looping subtree that runs until the shared done flag is

Selector

produced.

Notes

The preset shares the same state-key mapping across the LLM node and tool executor. For custom key layouts, pass a tailored ReActConfig.

源代码位于: jianmu/node/presets/react.py
def create_react_node(
    name: str = "ReActAgent",
    *,
    namespace: Optional[str] = None,
    model_client: ModelClient,
    tools: Optional[Seq[Tool]] = None,
    tool_providers: Optional[Seq[ToolProvider]] = None,
    tool_provider_context: Optional[Any] = None,
    tool_runner: Optional[Any] = None,
    constraints: Optional[Any] = None,
    context_builder: Optional[ContextBuilderProtocol] = None,
    config: Optional[ReActConfig] = None,
    tool_result_policy: Optional[ToolResultPolicy] = None,
    tool_result_policy_scope: Optional[ToolResultPolicyScope] = None,
) -> Selector:
    """Build a ReAct subtree with agent, tool, and completion nodes.

    This preset assembles the common ReAct loop:

    1. ask an agent LLM node for the next step
    2. execute emitted tool calls
    3. stop once ``done`` becomes truthy in state

    It is the quickest way to build a tool-using agent::

        root = create_react_node(
            model_client=my_model_client,
            tools=[CalculatorTool()],
        )

    Args:
        name: Root node name for the preset subtree.
        namespace: Optional state namespace for isolation.
        model_client: Model client used by the agent node.
        tools: Static tools available to the agent.
        tool_providers: Optional providers that contribute dynamic tools at
            assembly time.
        tool_provider_context: Optional context passed to tool providers during
            collection.
        tool_runner: Optional explicit tool runner override.
        constraints: Optional execution/guard constraints shared with the tool
            executor.
        context_builder: Optional prompt context builder passed to the agent
            LLM node.
        config: ReAct configuration. Defaults are used when omitted.
        tool_result_policy: Optional result policy evaluated after tool
            execution and before observation commit.
        tool_result_policy_scope: Optional predicate limiting the calls to
            which ``tool_result_policy`` applies. Without one, policy-enabled
            batches are conservatively serialized.

    Returns:
        Configured looping subtree that runs until the shared ``done`` flag is
        produced.

    Notes:
        The preset shares the same state-key mapping across the LLM node and
        tool executor. For custom key layouts, pass a tailored ``ReActConfig``.
    """
    tools = list(tools or [])
    provider_tools = ToolSet.from_providers(tool_providers or [], context=tool_provider_context)
    tools = provider_tools.with_tools(tools).tools
    config = config or ReActConfig()

    prefix = name
    keys = config.keys
    model_extra_params = dict(config.extra_params or {})
    legacy_tool_choice = model_extra_params.pop("tool_choice", None)
    initial_tool_choice = (
        config.tool_choice if config.tool_choice is not None else legacy_tool_choice
    )

    llm_keys = keys.model_copy(deep=True)
    llm_keys.text_output = keys.final_answer

    # Share the same state-key mapping across the LLM and tool executor.
    llm_config = AgentLLMConfig(
        model=config.model,
        temperature=config.temperature,
        max_tokens=config.max_tokens,
        top_p=config.top_p,
        top_k=config.top_k,
        timeout=config.timeout,
        stream=config.stream,
        max_budget_tokens=config.max_budget_tokens,
        resilience=config.resilience.model_copy(deep=True),
        system_prompt=config.system_prompt,
        extra_params=model_extra_params,
        keys=llm_keys,
    )

    tool_config = ToolExecutorConfig(keys=keys)

    # Precompute both structured tool schemas and text descriptions.
    toolset = ToolSet.from_tools(tools)
    tools_desc = toolset.describe()
    tools_schema = toolset.schemas()

    llm_node = _ReActAgentLLMNode(
        name=f"{prefix}/AgentLLM",
        model_client=model_client,
        config=llm_config,
        tools_schema=tools_schema,
        tools_description=tools_desc,
        namespace=namespace,
        context_builder=context_builder,
        initial_tool_choice=initial_tool_choice,
    )
    tool_executor = ToolExecutor(
        name=f"{prefix}/ToolExecutor",
        namespace=namespace,
        tools=tools,
        tool_runner=tool_runner,
        constraints=constraints,
        config=tool_config,
        tool_result_policy=tool_result_policy,
        tool_result_policy_scope=tool_result_policy_scope,
    )
    finalize_turn = _FinalizeAgentRoundNode(
        name=f"{prefix}/FinalizeTurn",
        namespace=namespace,
        actions_key=keys.actions,
        done_key=keys.done,
        rounds_key=keys.rounds,
    )

    loop_body = Sequence(
        name=f"{prefix}/ReActLoop",
        memory=True,
        children=[llm_node, tool_executor, finalize_turn]
    )

    loop_node = LoopUntilSuccess(
        name=f"{prefix}/NormalPath",
        max_iterations=config.max_iterations,
        child=loop_body
    )

    # Abort before the next iteration once the shared token budget is exceeded.
    def _check_abort_with_sm() -> bool:
        """Return whether policy completion or the shared budget should abort.

        Returns:
            ``True`` when a result policy completed with failure or the shared
            token budget is exhausted; otherwise ``False``.
        """
        if (
            tool_executor.status == Status.FAILURE
            and tool_executor.completion_status == Status.FAILURE
        ):
            return True
        sm = loop_node.state_manager
        if not sm or config.max_budget_tokens is None:
            return False

        current_usage = sm.get(keys.usage, namespace=namespace) or {}
        total_tokens = current_usage.get("total_tokens", 0)

        if total_tokens >= config.max_budget_tokens:
            logger.error(
                "🛑 [{}] Token budget exceeded ({} >= {}), aborting ReAct loop.",
                name, total_tokens, config.max_budget_tokens
            )
            write_runtime_termination(
                sm,
                build_budget_exceeded_termination(
                    node_name=name,
                    total_tokens=total_tokens,
                    max_budget_tokens=config.max_budget_tokens,
                ),
            )
            return True
        return False

    loop_node.abort_condition = _check_abort_with_sm
    if not config.soft_landing:
        return loop_node

    soft_landing_node = SoftLandingSummaryNode(
        name=f"{prefix}/SoftLandingFallback",
        namespace=namespace,
        model_client=model_client,
        context_builder=context_builder,
        config=AgentLLMConfig(
            model=config.model,
            temperature=config.temperature,
            max_tokens=config.max_tokens,
            top_p=config.top_p,
            top_k=config.top_k,
            timeout=config.timeout,
            stream=config.stream,
            max_budget_tokens=config.max_budget_tokens,
            resilience=config.resilience.model_copy(deep=True),
            system_prompt=config.system_prompt,
            extra_params=dict(model_extra_params),
            keys=keys,
        ),
        enabled=True,
    )
    return Selector(
        name=name,
        memory=True,
        children=[loop_node, soft_landing_node],
    )

create_plan_execute_node

create_plan_execute_node(
    name: str = "PlanExecute",
    *,
    namespace: Optional[str] = None,
    model_client: Optional[ModelClient] = None,
    tools: Optional[Sequence[Tool]] = None,
    tool_providers: Optional[Sequence[ToolProvider]] = None,
    tool_provider_context: Optional[Any] = None,
    tool_runner: Optional[Any] = None,
    constraints: Optional[Any] = None,
    config: Optional[PlanExecuteConfig] = None,
) -> LoopUntilSuccess

Build a plan-execute-review preset subtree.

This preset separates planning from execution. The planner writes an intermediate plan, the executor works against tools and message history, and an optional reviewer validates the result before the loop exits.

Use it when a task benefits from explicit decomposition rather than a single ReAct loop.

参数:

名称 类型 描述 默认
name str

Root node name for the preset subtree.

'PlanExecute'
namespace Optional[str]

Optional state namespace for isolation.

None
model_client Optional[ModelClient]

Model client used for planning, execution, and optional review. When omitted, ModelClient.resolve() is used.

None
tools Optional[Sequence[Tool]]

Static tools exposed to the executor.

None
tool_providers Optional[Sequence[ToolProvider]]

Optional providers that contribute dynamic tools at assembly time.

None
tool_provider_context Optional[Any]

Optional context passed to tool providers during collection.

None
tool_runner Optional[Any]

Optional explicit tool runner override for tool execution.

None
constraints Optional[Any]

Optional execution/guard constraints shared with the tool executor.

None
config Optional[PlanExecuteConfig]

Plan-execute configuration. Defaults are used when omitted.

None

返回:

类型 描述
LoopUntilSuccess

Configured looping subtree for planning, execution, and optional review.

Notes

Planner and executor use different logical state keys for their final outputs so the planner's plan does not overwrite the user-facing final answer.

源代码位于: jianmu/node/presets/plan_execute.py
def create_plan_execute_node(
    name: str = "PlanExecute",
    *,
    namespace: Optional[str] = None,
    model_client: Optional[ModelClient] = None,
    tools: Optional[Seq[Tool]] = None,
    tool_providers: Optional[Seq[ToolProvider]] = None,
    tool_provider_context: Optional[Any] = None,
    tool_runner: Optional[Any] = None,
    constraints: Optional[Any] = None,
    config: Optional[PlanExecuteConfig] = None,
) -> LoopUntilSuccess:
    """Build a plan-execute-review preset subtree.

    This preset separates planning from execution. The planner writes an
    intermediate plan, the executor works against tools and message history,
    and an optional reviewer validates the result before the loop exits.

    Use it when a task benefits from explicit decomposition rather than a
    single ReAct loop.

    Args:
        name: Root node name for the preset subtree.
        namespace: Optional state namespace for isolation.
        model_client: Model client used for planning, execution, and optional
            review. When omitted, ``ModelClient.resolve()`` is used.
        tools: Static tools exposed to the executor.
        tool_providers: Optional providers that contribute dynamic tools at
            assembly time.
        tool_provider_context: Optional context passed to tool providers during
            collection.
        tool_runner: Optional explicit tool runner override for tool execution.
        constraints: Optional execution/guard constraints shared with the tool
            executor.
        config: Plan-execute configuration. Defaults are used when omitted.

    Returns:
        Configured looping subtree for planning, execution, and optional review.

    Notes:
        Planner and executor use different logical state keys for their final
        outputs so the planner's plan does not overwrite the user-facing final
        answer.
    """
    static_tools = list(tools or [])
    provider_tools = ToolSet.from_providers(tool_providers or [], context=tool_provider_context)
    tools = provider_tools.with_tools(static_tools).tools
    config = config or PlanExecuteConfig()
    model_client = model_client or ModelClient.resolve()

    toolset = ToolSet.from_tools(tools)
    tools_desc = toolset.describe()
    tools_schema = toolset.schemas()
    planner_context_builder = ContextBuilder(
        providers=[
            StaticPromptProvider(config.plan_prompt if config.plan_prompt is not None else get_plan_execute_plan_prompt()),
            StateHistoryProvider(),
        ]
    )
    executor_context_builder = ContextBuilder(
        providers=[
            StaticPromptProvider(config.execute_prompt if config.execute_prompt is not None else get_plan_execute_execute_prompt()),
            *([ToolsDescProvider(tools_desc)] if tools_desc else []),
            StateHistoryProvider(),
        ]
    )

    keys = config.keys
    plan_field = str(config.plan_key or "plan")
    score_field = str(config.score_key or "score")

    # Keep the planner output separate from the executor's user-facing answer.
    planner_keys = StateKeys(
        messages=keys.messages,
        actions=keys.actions,
        done=keys.done,
        text_output=plan_field,
    )

    # Executor and tool execution share the same conversational state.
    executor_keys = StateKeys(
        messages=keys.messages,
        actions=keys.actions,
        done=keys.done,
        text_output=keys.final_answer,
        streaming_output=keys.streaming_output,
    )

    planner_config = AgentLLMConfig(
        model=config.model,
        temperature=config.temperature,
        max_tokens=config.max_tokens,
        top_p=config.top_p,
        top_k=config.top_k,
        timeout=config.timeout,
        system_prompt=config.plan_prompt,
        keys=planner_keys,
    )

    planner = AgentLLMNode(
        name=f"{name}/Planner",
        namespace=namespace,
        model_client=model_client,
        config=planner_config,
        tools_schema=None,
        tools_description="",
        context_builder=planner_context_builder,
    )

    executor_config = AgentLLMConfig(
        model=config.model,
        temperature=config.temperature,
        max_tokens=config.max_tokens,
        top_p=config.top_p,
        top_k=config.top_k,
        timeout=config.timeout,
        system_prompt=config.execute_prompt,
        stream=config.stream,
        keys=executor_keys,
    )

    executor = AgentLLMNode(
        name=f"{name}/Executor",
        namespace=namespace,
        model_client=model_client,
        config=executor_config,
        tools_schema=tools_schema,
        tools_description=tools_desc,
        context_builder=executor_context_builder,
    )

    nodes: list = [planner, executor]

    if tools:
        tool_config = ToolExecutorConfig(keys=executor_keys)
        tool_exec = ToolExecutor(
            name=f"{name}/ToolExecutor",
            namespace=namespace,
            tools=tools,
            tool_runner=tool_runner,
            constraints=constraints,
            config=tool_config,
        )
        nodes.append(tool_exec)

    def _review_passed() -> bool:
        if not config.enable_review:
            return True
        score = loop_body.state_manager.get(score_field, namespace=namespace, default=None) if loop_body.state_manager else None
        try:
            return score is not None and float(score) >= float(config.review_threshold)
        except (TypeError, ValueError):
            return False

    if config.enable_review:
        reviewer = EvaluationNode(
            name=f"{name}/Reviewer",
            namespace=namespace,
            model_client=model_client,
            config=AgentLLMConfig(
                model=config.model,
                temperature=config.temperature,
                max_tokens=config.max_tokens,
                top_p=config.top_p,
                top_k=config.top_k,
                timeout=config.timeout,
                system_prompt=config.review_prompt or get_plan_execute_review_prompt(),
            ),
            output_keys=EvaluationStateKeys(score=score_field),
        )
        nodes.append(reviewer)

    nodes.append(
        _FinalizeAgentRoundNode(
            name=f"{name}/FinalizeRound",
            namespace=namespace,
            actions_key=keys.actions,
            done_key=keys.done,
            rounds_key=keys.rounds,
            extra_completion_check=_review_passed,
        )
    )

    loop_body = PySequence(name=f"{name}/Loop", memory=True, children=nodes)
    return LoopUntilSuccess(name=name, max_iterations=config.max_rounds, child=loop_body)