跳转至

jianmu.engine

适用对象:运行时扩展开发者 / 核心维护者 / 高阶使用者 是否必读:按需 相关模块:jianmu.node, jianmu.message, jianmu.memory

1. 模块职责

jianmu.engine 是运行时核心层,负责行为树异步执行、状态管理和运行上下文注入。

大多数普通用户会通过 jianmu 根入口或 jianmu.node 间接使用这里的能力;只有在你需要定制运行时行为时,才需要直接进入这个模块。

2. 适合查什么

  • AsyncBehaviour
  • StateManager
  • ReactiveRunner
  • RunContext

3. 使用建议

  • 自定义节点时,通常依赖这里提供的运行时注入语义
  • 自定义调度或执行模型时,再直接使用 ReactiveRunner
  • LoopUntilSuccess 这类 Jianmu 自定义组合节点优先从 jianmu.node 查看

4. 最小示例

from jianmu.engine import ReactiveRunner, RunContext, StateManager

state_manager = StateManager(MyState)
state_manager.initialize()

runner = ReactiveRunner(
    root=my_tree,
    state_manager=state_manager,
    ctx=RunContext(),
)

5. 常见入口

  • 想管理 durable state:看 StateManager
  • 想驱动树运行:看 ReactiveRunner
  • 想传运行时依赖:看 RunContext
  • 想用 Jianmu 自定义组合节点:看 jianmu.node.LoopUntilSuccess

6. 注意事项

  • 普通应用通常会从 jianmu 根入口或 jianmu.node 间接使用这一层
  • 只有在你要自定义运行时行为时,才需要直接长期依赖 jianmu.engine
  • AsyncBehaviour 更偏运行时基类;面向用户的同步/异步节点概念仍优先看 jianmu.Node / jianmu.AsyncNode

7. API 参考

Agent

Agent

Agent(
    root: Behaviour,
    *,
    state_schema: type[BaseModel] | None = None,
    state_manager: StateManager | None = None,
    context: RunContext | None = None,
    setup_timeout: float = 15.0,
    checkpointer: CheckpointerProtocol | None = None,
    thread_id: str | None = None,
    checkpoint_interval: int | None = None,
    max_fps: float | None = None,
)

User-facing single-agent facade that wraps ReactiveRunner.

Agent is a lightweight shell around a prepared behaviour tree, its state manager, runtime context, and one configured ReactiveRunner. It does not replace the runner; it centralizes high-level construction, default runtime configuration, and checkpoint wiring.

属性:

名称 类型 描述
root

Root behaviour tree bound to the facade.

state_manager

Backing state manager used for all agent execution.

context

Runtime dependency context shared with the wrapped runner.

runner ReactiveRunner

Wrapped low-level ReactiveRunner instance.

Create an Agent facade from a prepared behaviour tree.

参数:

名称 类型 描述 默认
root Behaviour

Root behaviour tree node.

必需
state_schema type[BaseModel] | None

Optional state schema used when creating a state manager implicitly.

None
state_manager StateManager | None

Optional prebuilt state manager.

None
context RunContext | None

Optional runtime dependency context.

None
setup_timeout float

Maximum setup time passed to the wrapped runner.

15.0
checkpointer CheckpointerProtocol | None

Optional explicit default checkpointer.

None
thread_id str | None

Optional default checkpoint thread identifier.

None
checkpoint_interval int | None

Optional default checkpoint save interval.

None
max_fps float | None

Optional default runner maximum FPS.

None

引发:

类型 描述
ValueError

If neither state_manager nor state_schema is provided.

源代码位于: jianmu/engine/agent.py
def __init__(
    self,
    root: Behaviour,
    *,
    state_schema: type[BaseModel] | None = None,
    state_manager: StateManager | None = None,
    context: RunContext | None = None,
    setup_timeout: float = 15.0,
    checkpointer: CheckpointerProtocol | None = None,
    thread_id: str | None = None,
    checkpoint_interval: int | None = None,
    max_fps: float | None = None,
):
    """Create an ``Agent`` facade from a prepared behaviour tree.

    Args:
        root: Root behaviour tree node.
        state_schema: Optional state schema used when creating a state
            manager implicitly.
        state_manager: Optional prebuilt state manager.
        context: Optional runtime dependency context.
        setup_timeout: Maximum setup time passed to the wrapped runner.
        checkpointer: Optional explicit default checkpointer.
        thread_id: Optional default checkpoint thread identifier.
        checkpoint_interval: Optional default checkpoint save interval.
        max_fps: Optional default runner maximum FPS.

    Raises:
        ValueError: If neither ``state_manager`` nor ``state_schema`` is provided.
    """
    if state_manager is None:
        if state_schema is None:
            raise ValueError("Agent requires either 'state_manager' or 'state_schema'")
        state_manager = StateManager(schema=state_schema)
        state_manager.initialize()

    ctx = context or RunContext()
    resolved_checkpointer = (
        checkpointer if checkpointer is not None else _resolve_default_checkpointer()
    )
    defaults = _resolve_agent_runtime_defaults(
        thread_id=thread_id,
        checkpoint_interval=checkpoint_interval,
        max_fps=max_fps,
    )
    runner = ReactiveRunner(
        root,
        state_manager,
        ctx=ctx,
        setup_timeout=setup_timeout,
        max_pending_wakeups=None,
        hot_loop_warn_factor=None,
    )

    self.root = root
    self.state_manager = state_manager
    self.context = ctx
    self._runner = runner
    self._checkpointer = resolved_checkpointer
    self._defaults = defaults

runner property

runner: ReactiveRunner

Return the wrapped low-level reactive runner.

返回:

类型 描述
ReactiveRunner

The resulting ReactiveRunner value.

state property

state: Any

Return the current state snapshot from the state manager.

返回:

类型 描述
Any

The resulting Any value.

run async

run(
    input_data: dict[str, Any] | None = None,
    *,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    thread_id: str | None = None,
    max_fps: float | None = None,
) -> Any

Run the wrapped agent to completion.

参数:

名称 类型 描述 默认
input_data dict[str, Any] | None

Optional initial input state merged before the run.

None
reset_tree bool

Whether to reset the tree before execution.

True
reset_data bool

Whether to reset the backing state manager first.

False
max_ticks int | None

Optional hard tick limit.

None
timeout_s float | None

Optional wall-clock timeout for the full run.

None
checkpointer CheckpointerProtocol | None

Optional explicit run-time checkpointer override.

None
checkpoint_interval int | None

Optional explicit checkpoint interval override.

None
event_driven_checkpoint bool

Whether to save checkpoints only when checkpointable state/tree status changes.

False
thread_id str | None

Optional explicit checkpoint thread identifier override.

None
max_fps float | None

Optional explicit runner FPS cap override.

None

返回:

类型 描述
Any

Final root-node status from the wrapped runner.

源代码位于: jianmu/engine/agent.py
async def run(
    self,
    input_data: dict[str, Any] | None = None,
    *,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    thread_id: str | None = None,
    max_fps: float | None = None,
) -> Any:
    """Run the wrapped agent to completion.

    Args:
        input_data: Optional initial input state merged before the run.
        reset_tree: Whether to reset the tree before execution.
        reset_data: Whether to reset the backing state manager first.
        max_ticks: Optional hard tick limit.
        timeout_s: Optional wall-clock timeout for the full run.
        checkpointer: Optional explicit run-time checkpointer override.
        checkpoint_interval: Optional explicit checkpoint interval override.
        event_driven_checkpoint: Whether to save checkpoints only when
            checkpointable state/tree status changes.
        thread_id: Optional explicit checkpoint thread identifier override.
        max_fps: Optional explicit runner FPS cap override.

    Returns:
        Final root-node status from the wrapped runner.
    """
    return await self._runner.run(
        input_data=input_data,
        reset_tree=reset_tree,
        reset_data=reset_data,
        max_ticks=max_ticks,
        timeout_s=timeout_s,
        checkpointer=checkpointer if checkpointer is not None else self._checkpointer,
        checkpoint_interval=_resolve_checkpoint_interval_override(
            default=self._defaults.checkpoint_interval,
            checkpoint_interval=checkpoint_interval,
            event_driven_checkpoint=event_driven_checkpoint,
        ),
        thread_id=thread_id or self._defaults.thread_id,
        max_fps=max_fps if max_fps is not None else self._defaults.max_fps,
    )

step async

step(
    obs: dict[str, Any] | None = None,
    *,
    yield_to_async: bool = False,
) -> dict[str, Any]

Execute one runner step in step mode.

参数:

名称 类型 描述 默认
obs dict[str, Any] | None

Optional observation payload merged before the step.

None
yield_to_async bool

Whether to yield once after the tick.

False

返回:

类型 描述
dict[str, Any]

Step-scoped state fields returned by the wrapped runner.

源代码位于: jianmu/engine/agent.py
async def step(
    self,
    obs: dict[str, Any] | None = None,
    *,
    yield_to_async: bool = False,
) -> dict[str, Any]:
    """Execute one runner step in step mode.

    Args:
        obs: Optional observation payload merged before the step.
        yield_to_async: Whether to yield once after the tick.

    Returns:
        Step-scoped state fields returned by the wrapped runner.
    """
    return await self._runner.step(obs=obs, yield_to_async=yield_to_async)

run_until_suspend async

run_until_suspend(
    input_data: dict[str, Any] | None = None,
    *,
    resume_data: dict[str, Any] | None = None,
    resume_request_id: str | None = None,
    restore: str | RestorePolicy = RestorePolicy.NEVER,
    suspension_mode: str
    | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    thread_id: str | None = None,
    max_fps: float | None = None,
) -> RunResult

Run until terminal completion or a publishable suspension point.

参数:

名称 类型 描述 默认
input_data dict[str, Any] | None

The input_data value.

None
resume_data dict[str, Any] | None

The resume_data value.

None
resume_request_id str | None

Active suspension id required when resume_data is provided.

None
restore str | RestorePolicy

The restore value.

NEVER
suspension_mode str | SuspensionMode

The suspension_mode value.

YIELD
reset_tree bool

The reset_tree value.

True
reset_data bool

The reset_data value.

False
max_ticks int | None

Collection of max tick values.

None
timeout_s float | None

Collection of timeout values.

None
checkpointer CheckpointerProtocol | None

The checkpointer value.

None
checkpoint_interval int | None

The checkpoint_interval value.

None
event_driven_checkpoint bool

Whether to save checkpoints only when checkpointable state/tree status changes.

False
thread_id str | None

Identifier for thread.

None
max_fps float | None

Collection of max fp values.

None

返回:

类型 描述
RunResult

The resulting RunResult value.

源代码位于: jianmu/engine/agent.py
async def run_until_suspend(
    self,
    input_data: dict[str, Any] | None = None,
    *,
    resume_data: dict[str, Any] | None = None,
    resume_request_id: str | None = None,
    restore: str | RestorePolicy = RestorePolicy.NEVER,
    suspension_mode: str | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpointer: CheckpointerProtocol | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    thread_id: str | None = None,
    max_fps: float | None = None,
) -> RunResult:
    """Run until terminal completion or a publishable suspension point.

    Args:
        input_data: The `input_data` value.
        resume_data: The `resume_data` value.
        resume_request_id: Active suspension id required when resume_data is provided.
        restore: The `restore` value.
        suspension_mode: The `suspension_mode` value.
        reset_tree: The `reset_tree` value.
        reset_data: The `reset_data` value.
        max_ticks: Collection of max tick values.
        timeout_s: Collection of timeout  values.
        checkpointer: The `checkpointer` value.
        checkpoint_interval: The `checkpoint_interval` value.
        event_driven_checkpoint: Whether to save checkpoints only when
            checkpointable state/tree status changes.
        thread_id: Identifier for thread.
        max_fps: Collection of max fp values.

    Returns:
        The resulting `RunResult` value.
    """
    return await self._runner.run_until_suspend(
        input_data=input_data,
        resume_data=resume_data,
        resume_request_id=resume_request_id,
        restore=restore,
        suspension_mode=suspension_mode,
        reset_tree=reset_tree,
        reset_data=reset_data,
        max_ticks=max_ticks,
        timeout_s=timeout_s,
        checkpointer=checkpointer if checkpointer is not None else self._checkpointer,
        checkpoint_interval=_resolve_checkpoint_interval_override(
            default=self._defaults.checkpoint_interval,
            checkpoint_interval=checkpoint_interval,
            event_driven_checkpoint=event_driven_checkpoint,
        ),
        thread_id=thread_id or self._defaults.thread_id,
        max_fps=max_fps if max_fps is not None else self._defaults.max_fps,
    )

resume_interaction async

resume_interaction(
    *,
    thread_id: str | None = None,
    request_id: str,
    payload: dict[str, Any],
    checkpointer: CheckpointerProtocol | None = None,
    suspension_mode: str
    | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    max_fps: float | None = None,
) -> ResumeInteractionResult

Resume a suspended interaction and return a facade-friendly result.

This facade wraps ReactiveRunner.resume_and_continue and preserves its strict request_id semantics. Approval suspensions resume through the shared RunContext.approval_manager bridge; when the bridge prerequisites are missing, the method returns a structured recoverable resume result instead of raising to the caller.

参数:

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

Optional checkpoint thread identifier override.

None
request_id str

Active suspension request identifier to resume.

必需
payload dict[str, Any]

Host-provided resume payload.

必需
checkpointer CheckpointerProtocol | None

Optional explicit run-time checkpointer override.

None
suspension_mode str | SuspensionMode

Host-facing handling mode for any subsequent suspension.

YIELD
reset_tree bool

Whether to reset tree execution before resuming.

True
reset_data bool

Whether to reset the backing state manager first.

False
max_ticks int | None

Optional hard tick limit.

None
timeout_s float | None

Optional wall-clock timeout for the resumed run.

None
checkpoint_interval int | None

Optional checkpoint interval override.

None
event_driven_checkpoint bool

Whether to save checkpoints only when checkpointable state/tree status changes.

False
max_fps float | None

Optional runner FPS cap override.

None

返回:

类型 描述
ResumeInteractionResult

Structured resume status plus the run result when execution starts.

源代码位于: jianmu/engine/agent.py
async def resume_interaction(
    self,
    *,
    thread_id: str | None = None,
    request_id: str,
    payload: dict[str, Any],
    checkpointer: CheckpointerProtocol | None = None,
    suspension_mode: str | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int | None = None,
    timeout_s: float | None = None,
    checkpoint_interval: int | None = None,
    event_driven_checkpoint: bool = False,
    max_fps: float | None = None,
) -> ResumeInteractionResult:
    """Resume a suspended interaction and return a facade-friendly result.

    This facade wraps ``ReactiveRunner.resume_and_continue`` and preserves
    its strict ``request_id`` semantics. Approval suspensions resume
    through the shared ``RunContext.approval_manager`` bridge; when the
    bridge prerequisites are missing, the method returns a structured
    recoverable resume result instead of raising to the caller.

    Args:
        thread_id: Optional checkpoint thread identifier override.
        request_id: Active suspension request identifier to resume.
        payload: Host-provided resume payload.
        checkpointer: Optional explicit run-time checkpointer override.
        suspension_mode: Host-facing handling mode for any subsequent suspension.
        reset_tree: Whether to reset tree execution before resuming.
        reset_data: Whether to reset the backing state manager first.
        max_ticks: Optional hard tick limit.
        timeout_s: Optional wall-clock timeout for the resumed run.
        checkpoint_interval: Optional checkpoint interval override.
        event_driven_checkpoint: Whether to save checkpoints only when
            checkpointable state/tree status changes.
        max_fps: Optional runner FPS cap override.

    Returns:
        Structured resume status plus the run result when execution starts.
    """
    try:
        run_result = await self._runner.resume_and_continue(
            thread_id=thread_id or self._defaults.thread_id,
            request_id=request_id,
            payload=payload,
            checkpointer=checkpointer if checkpointer is not None else self._checkpointer,
            suspension_mode=suspension_mode,
            reset_tree=reset_tree,
            reset_data=reset_data,
            max_ticks=max_ticks,
            timeout_s=timeout_s,
            checkpoint_interval=_resolve_checkpoint_interval_override(
                default=self._defaults.checkpoint_interval,
                checkpoint_interval=checkpoint_interval,
                event_driven_checkpoint=event_driven_checkpoint,
            ),
            max_fps=max_fps if max_fps is not None else self._defaults.max_fps,
        )
    except ResumeError as exc:
        return ResumeInteractionResult(resume=exc.result, run=None)
    return ResumeInteractionResult(
        resume=ResumeSuspensionResult(
            status=ResumeStatus.SUCCESS,
            request_id=request_id,
            message="Resume request accepted.",
            payload=dict(payload or {}),
            retryable=False,
        ),
        run=run_result,
    )

resume async

resume(**kwargs: Any) -> Any

Continue agent execution without resetting tree or state.

This is a convenience alias for continuing the current in-memory runner state. It is not a durable checkpoint restore API.

参数:

名称 类型 描述 默认
**kwargs Any

Additional keyword arguments forwarded to run().

{}

返回:

类型 描述
Any

Final root-node status from the wrapped runner.

源代码位于: jianmu/engine/agent.py
async def resume(self, **kwargs: Any) -> Any:
    """Continue agent execution without resetting tree or state.

    This is a convenience alias for continuing the current in-memory
    runner state. It is not a durable checkpoint restore API.

    Args:
        **kwargs: Additional keyword arguments forwarded to ``run()``.

    Returns:
        Final root-node status from the wrapped runner.
    """
    return await self.run(reset_tree=False, reset_data=False, **kwargs)

reset

reset(reset_data: bool = True) -> None

Reset the wrapped runner.

参数:

名称 类型 描述 默认
reset_data bool

Whether to reinitialize the state manager.

True
源代码位于: jianmu/engine/agent.py
def reset(self, reset_data: bool = True) -> None:
    """Reset the wrapped runner.

    Args:
        reset_data: Whether to reinitialize the state manager.
    """
    self._runner.reset(reset_data=reset_data)

AsyncBehaviour

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()")

StateManager

StateManager

StateManager(schema: Type[T])

Thread-safe state container with validation, reducers, and change notifications.

StateManager is the durable state surface for a Jianmu workflow. It owns one Pydantic schema instance, validates all writes, supports reducer-style merges, and notifies listeners such as ReactiveRunner when state changes should wake the tree.

Typical usage looks like::

from typing import Annotated
from pydantic import BaseModel

class AgentState(BaseModel):
    task: str = ""
    messages: list = []
    scratchpad: Annotated[list[str], Ephemeral("step")] = []

sm = StateManager(AgentState)
sm.initialize({"task": "summarize this repository"})
sm.update({"messages": ["hello"]})
current = sm.get()

Use the root schema for durable state, and namespace=... only for runtime-scoped extra fields when the schema allows extra='allow'.

Create a state manager for a Pydantic schema.

参数:

名称 类型 描述 默认
schema Type[T]

Pydantic model type that defines the durable state surface.

必需
源代码位于: jianmu/engine/state.py
def __init__(self, schema: Type[T]):
    """Create a state manager for a Pydantic schema.

    Args:
        schema: Pydantic model type that defines the durable state surface.
    """
    self.schema = schema
    self.reducers: Dict[str, Callable[[Any, Any], Any]] = {}
    # name -> (default_value, default_factory, scope)
    self._ephemeral_fields: Dict[str, tuple] = {}

    self._listeners: List[Callable[[], None]] = []
    self._listeners_lock = threading.Lock()
    self._debug_snapshots: List[Dict[str, Any]] = []
    self._runtime_metadata: Dict[str, Any] = {}
    self._revision = 0

    self._lock = threading.Lock()
    self._data: Optional[T] = None

    # Per-field locks let multiple store wrappers serialize read-modify-write
    # operations against the same logical key.
    self._field_write_locks: Dict[Any, threading.Lock] = {}
    self._field_write_locks_lock = threading.Lock()

    self._parse_schema()

revision property

revision: int

Return a monotonic durable-state revision for checkpoint cursors.

This is runner bookkeeping, not a business-level data version.

subscribe

subscribe(callback: Callable[[], None])

Register a state-change listener.

参数:

名称 类型 描述 默认
callback Callable[[], None]

Zero-argument callable invoked after successful writes or explicit signal notifications.

必需
源代码位于: jianmu/engine/state.py
def subscribe(self, callback: Callable[[], None]):
    """Register a state-change listener.

    Args:
        callback: Zero-argument callable invoked after successful writes or
            explicit signal notifications.
    """
    with self._listeners_lock:
        self._listeners.append(callback)

unsubscribe

unsubscribe(callback: Callable[[], None])

Remove a previously registered listener.

参数:

名称 类型 描述 默认
callback Callable[[], None]

Listener previously passed to :meth:subscribe.

必需
源代码位于: jianmu/engine/state.py
def unsubscribe(self, callback: Callable[[], None]):
    """Remove a previously registered listener.

    Args:
        callback: Listener previously passed to :meth:`subscribe`.
    """
    with self._listeners_lock:
        try:
            self._listeners.remove(callback)
        except ValueError:
            pass

signal

signal()

Notify listeners without mutating state.

源代码位于: jianmu/engine/state.py
def signal(self):
    """Notify listeners without mutating state."""
    if not state_write_allowed():
        _guarded_write_blocked("signal")
        return
    self._notify_listeners()

initialize

initialize(initial_state: Optional[Dict[str, Any]] = None)

Initialize and validate the state model.

参数:

名称 类型 描述 默认
initial_state Optional[Dict[str, Any]]

Optional initial values used to build the backing schema instance.

None

引发:

类型 描述
ValueError

If the supplied state does not satisfy the schema.

源代码位于: jianmu/engine/state.py
def initialize(self, initial_state: Optional[Dict[str, Any]] = None):
    """Initialize and validate the state model.

    Args:
        initial_state: Optional initial values used to build the backing schema instance.

    Raises:
        ValueError: If the supplied state does not satisfy the schema.
    """
    data = initial_state or {}
    try:
        self._data = self.schema(**data)
    except ValidationError as e:
        raise ValueError(f"❌ [StateManager] Init Error: {e}")
    with self._lock:
        self._bump_revision()

get

get(
    key: str | None = None,
    namespace: str | None = None,
    default: Any = None,
) -> T | Any

Return the whole state model or a single field lookup.

参数:

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

Optional field name to read. When omitted, returns a validated copy of the full state model.

None
namespace str | None

Optional runtime namespace prefix for extra fields.

None
default Any

Fallback value used when the requested field is missing.

None

返回:

类型 描述
T | Any

The full state model copy or the resolved field value.

引发:

类型 描述
ValueError

If the full model is requested before initialization and the schema cannot be constructed from defaults alone.

源代码位于: jianmu/engine/state.py
def get(self, key: str | None = None, namespace: str | None = None, default: Any = None) -> T | Any:
    """Return the whole state model or a single field lookup.

    Args:
        key: Optional field name to read. When omitted, returns a validated copy
            of the full state model.
        namespace: Optional runtime namespace prefix for extra fields.
        default: Fallback value used when the requested field is missing.

    Returns:
        The full state model copy or the resolved field value.

    Raises:
        ValueError: If the full model is requested before initialization and the
            schema cannot be constructed from defaults alone.
    """
    full_key = f"{self._normalize_namespace(namespace)}{key}" if key is not None else None

    with self._lock:
        if self._data is None:
            if full_key is not None:
                return default
            try:
                return self.schema()
            except ValidationError as e:
                raise ValueError(
                    "❌ [StateManager] State not initialized. "
                    "Call initialize() with required fields or update() with required values first."
                ) from e

        if full_key is not None:
            # Namespaced runtime data lives in model extras instead of declared schema fields.
            extra = getattr(self._data, "model_extra", None)
            if extra is not None and full_key in extra:
                val = extra[full_key]
                return val if val is not None else default

            if full_key in self.schema.model_fields:
                val = getattr(self._data, full_key)
                return val if val is not None else default

            return default

        return self.schema.model_validate(self._data.model_dump())

update

update(
    updates: Dict[str, Any],
    signal: bool = True,
    namespace: str | None = None,
)

Apply validated updates, then notify listeners if requested.

参数:

名称 类型 描述 默认
updates Dict[str, Any]

Field values to merge into the current state.

必需
signal bool

Whether to notify listeners after the write succeeds.

True
namespace str | None

Optional runtime namespace for extra fields.

None

引发:

类型 描述
ValueError

If the write targets invalid schema fields or fails validation.

RuntimeError

If a registered reducer raises during the merge.

源代码位于: jianmu/engine/state.py
def update(self, updates: Dict[str, Any], signal: bool = True, namespace: str | None = None):
    """Apply validated updates, then notify listeners if requested.

    Args:
        updates: Field values to merge into the current state.
        signal: Whether to notify listeners after the write succeeds.
        namespace: Optional runtime namespace for extra fields.

    Raises:
        ValueError: If the write targets invalid schema fields or fails validation.
        RuntimeError: If a registered reducer raises during the merge.
    """
    if not state_write_allowed():
        _guarded_write_blocked("update")
        return
    prefix = self._normalize_namespace(namespace)

    # Namespace writes are reserved for runtime-only extras. They must never
    # overlap with declared schema fields, otherwise callers can silently
    # bypass reducer and validation expectations on the root model surface.
    if prefix and updates:
        prefixed_keys = [f"{prefix}{k}" for k in updates.keys()]

        collisions = [k for k in prefixed_keys if k in self.schema.model_fields]
        if collisions:
            raise ValueError(
                f"❌ [StateManager] Namespace writes cannot target schema fields: {collisions}. "
                "Use non-namespaced update() for declared state fields."
            )
        allows_extra = self.schema.model_config.get("extra") == "allow"
        if not allows_extra:
            undefined_keys = [k for k in prefixed_keys if k not in self.schema.model_fields]
            if undefined_keys:
                raise ValueError(
                    f"❌ [StateManager] Schema extra='allow' is required! "
                    f"Attempted to dynamically assign prefixed scoped keys {undefined_keys} "
                    f"but model schema '{self.schema.__name__}' is strictly typed."
                )
    if prefix:
        updates = {f"{prefix}{k}": v for k, v in updates.items()}

    with self._lock:
        current_data = self._data.model_dump() if self._data is not None else {}
        pending_writes = {}

        for name, update_val in updates.items():
            if name in self.reducers:
                reducer = self.reducers[name]
                old_val = current_data.get(name)
                if old_val is None and self._data is None:
                    field_info = self.schema.model_fields.get(name)
                    if field_info is not None:
                        try:
                            if not field_info.is_required():
                                old_val = field_info.get_default(call_default_factory=True)
                        except Exception:
                            pass
                try:
                    final_val = reducer(old_val, update_val)
                except Exception as e:
                    raise RuntimeError(f"❌ [StateManager] Reducer '{name}' failed: {e}")
            else:
                final_val = update_val

            pending_writes[name] = final_val

        merged_data = current_data.copy()
        merged_data.update(pending_writes)

        if merged_data != current_data:
            try:
                self._data = self.schema(**merged_data)
            except ValidationError as e:
                raise ValueError(
                    f"❌ [StateManager] Update Validation Failed: {e}. "
                    "If this is the first update, ensure all required fields are provided "
                    "or call initialize() first."
                )
            self._bump_revision()

    if signal:
        self._notify_listeners()

clear

clear(namespace: str, signal: bool = True)

Delete runtime-only fields under a namespace prefix.

参数:

名称 类型 描述 默认
namespace str

Namespace whose extra fields should be removed.

必需
signal bool

Whether to notify listeners after clearing data.

True
源代码位于: jianmu/engine/state.py
def clear(self, namespace: str, signal: bool = True):
    """Delete runtime-only fields under a namespace prefix.

    Args:
        namespace: Namespace whose extra fields should be removed.
        signal: Whether to notify listeners after clearing data.
    """
    if not state_write_allowed():
        _guarded_write_blocked("clear")
        return
    prefix = self._normalize_namespace(namespace)
    if not prefix:
        return

    with self._lock:
        if self._data is None:
            return

        extra = getattr(self._data, "model_extra", None)
        if not extra:
            return

        keys_to_delete = [k for k in list(extra.keys()) if k.startswith(prefix)]
        if not keys_to_delete:
            return

        for k in keys_to_delete:
            extra.pop(k, None)
        self._bump_revision()

    if signal:
        self.signal()

snapshot_namespace

snapshot_namespace(
    namespace: str, strip_prefix: bool = True
) -> Dict[str, Any]

Return a deep-copied snapshot of runtime namespace data.

参数:

名称 类型 描述 默认
namespace str

Namespace to snapshot.

必需
strip_prefix bool

Whether returned keys should omit the namespace prefix.

True

返回:

类型 描述
Dict[str, Any]

A deep-copied mapping of namespace-local state values.

源代码位于: jianmu/engine/state.py
def snapshot_namespace(self, namespace: str, strip_prefix: bool = True) -> Dict[str, Any]:
    """Return a deep-copied snapshot of runtime namespace data.

    Args:
        namespace: Namespace to snapshot.
        strip_prefix: Whether returned keys should omit the namespace prefix.

    Returns:
        A deep-copied mapping of namespace-local state values.
    """
    prefix = self._normalize_namespace(namespace)
    if not prefix:
        return {}

    with self._lock:
        if self._data is None:
            return {}

        extra = getattr(self._data, "model_extra", None) or {}
        snap: Dict[str, Any] = {}
        for key, value in extra.items():
            if key.startswith(prefix):
                out_key = key[len(prefix):] if strip_prefix else key
                snap[out_key] = copy.deepcopy(value)
        return snap

merge_dict

merge_dict(field: str, updates: Dict[str, Any])

Atomically merge a dictionary field.

参数:

名称 类型 描述 默认
field str

Dictionary-typed schema field to merge into.

必需
updates Dict[str, Any]

Partial mapping applied with last-write-wins semantics.

必需

引发:

类型 描述
ValueError

If the merged dictionary fails schema validation.

源代码位于: jianmu/engine/state.py
def merge_dict(self, field: str, updates: Dict[str, Any]):
    """Atomically merge a dictionary field.

    Args:
        field: Dictionary-typed schema field to merge into.
        updates: Partial mapping applied with last-write-wins semantics.

    Raises:
        ValueError: If the merged dictionary fails schema validation.
    """
    if not state_write_allowed():
        _guarded_write_blocked("merge_dict")
        return
    with self._lock:
        if self._data is None:
            return
        current_data = self._data.model_dump()
        current = current_data.get(field, getattr(self._data, field, {})) or {}
        if not isinstance(current, dict):
            current = {}
        merged = {**current, **updates}
        merged_data = current_data.copy()
        merged_data[field] = merged
        if merged_data != current_data:
            try:
                self._data = self.schema(**merged_data)
            except ValidationError as e:
                raise ValueError(f"❌ [StateManager] merge_dict Validation Failed: {e}")
            self._bump_revision()
    self.signal()

transform_field

transform_field(
    key: str,
    transform: Callable[[Any], Any],
    *,
    namespace: str | None = None,
    signal: bool = True,
) -> Any

Atomically transform one field and persist the validated result.

This is the safest way to perform read-modify-write updates for mutable fields such as message histories or counters::

sm.transform_field("messages", lambda current: list(current or []) + [msg])

参数:

名称 类型 描述 默认
key str

Logical field name to transform.

必需
transform Callable[[Any], Any]

Pure function receiving the current field value.

必需
namespace str | None

Optional runtime namespace prefix.

None
signal bool

Whether to notify listeners after a successful write.

True

返回:

类型 描述
Any

The transformed field value returned by transform.

引发:

类型 描述
ValueError

If namespace writes target schema fields, if extra='allow' is not set, or if validation fails.

源代码位于: jianmu/engine/state.py
def transform_field(
    self,
    key: str,
    transform: Callable[[Any], Any],
    *,
    namespace: str | None = None,
    signal: bool = True,
) -> Any:
    """Atomically transform one field and persist the validated result.

    This is the safest way to perform read-modify-write updates for mutable
    fields such as message histories or counters::

        sm.transform_field("messages", lambda current: list(current or []) + [msg])

    Args:
        key: Logical field name to transform.
        transform: Pure function receiving the current field value.
        namespace: Optional runtime namespace prefix.
        signal: Whether to notify listeners after a successful write.

    Returns:
        The transformed field value returned by ``transform``.

    Raises:
        ValueError: If namespace writes target schema fields, if extra='allow' is not set, or if validation fails.
    """
    if not state_write_allowed():
        _guarded_write_blocked("transform_field")
        return self.get(key, namespace=namespace, default=None)
    prefix = self._normalize_namespace(namespace)
    full_key = f"{prefix}{key}" if prefix else key

    if prefix:
        if full_key in self.schema.model_fields:
            raise ValueError(
                f"❌ [StateManager] Namespace writes cannot target schema fields: {[full_key]}. "
                "Use non-namespaced transform_field() for declared state fields."
            )
        allows_extra = self.schema.model_config.get("extra") == "allow"
        if not allows_extra:
            raise ValueError(
                f"❌ [StateManager] Schema extra='allow' is required for namespaced writes ({full_key})."
            )

    with self._get_write_lock(namespace, key):
        with self._lock:
            current_data = self._data.model_dump() if self._data is not None else {}
            if self._data is None:
                current_value = None
            elif prefix:
                current_value = current_data.get(full_key)
            elif full_key in self.schema.model_fields:
                current_value = getattr(self._data, full_key, None)
            else:
                current_value = current_data.get(full_key)
            next_value = transform(copy.deepcopy(current_value))
            if current_value != next_value or full_key not in current_data:
                current_data[full_key] = next_value
                try:
                    self._data = self.schema(**current_data)
                except ValidationError as e:
                    raise ValueError(f"❌ [StateManager] transform_field Validation Failed: {e}")
                self._bump_revision()

    if signal:
        self.signal()

    return next_value

dump_checkpoint

dump_checkpoint() -> Dict[str, Any]

Export the full state payload for checkpoint persistence.

返回:

类型 描述
Dict[str, Any]

Serializable checkpoint payload containing the current state.

源代码位于: jianmu/engine/state.py
def dump_checkpoint(self) -> Dict[str, Any]:
    """Export the full state payload for checkpoint persistence.

    Returns:
        Serializable checkpoint payload containing the current state.
    """
    with self._lock:
        state_data = self._data.model_dump() if self._data else {}
        return {
            "state": state_data,
            "runtime_metadata": copy.deepcopy(self._runtime_metadata),
        }

set_runtime_metadata

set_runtime_metadata(key: str, value: Any) -> None

Persist runner-internal metadata outside the user state schema.

参数:

名称 类型 描述 默认
key str

The key value.

必需
value Any

Value to apply.

必需
源代码位于: jianmu/engine/state.py
def set_runtime_metadata(self, key: str, value: Any) -> None:
    """Persist runner-internal metadata outside the user state schema.

    Args:
        key: The `key` value.
        value: Value to apply.
    """
    if not state_write_allowed():
        _guarded_write_blocked("set_runtime_metadata")
        return
    with self._lock:
        normalized_key = str(key)
        next_value = copy.deepcopy(value)
        if self._runtime_metadata.get(normalized_key) == next_value:
            return
        self._runtime_metadata[normalized_key] = next_value
        self._bump_revision()
    self.signal()

get_runtime_metadata

get_runtime_metadata(key: str, default: Any = None) -> Any

Return one runner-internal metadata value.

参数:

名称 类型 描述 默认
key str

The key value.

必需
default Any

Default value to use when no explicit value is available.

None

返回:

类型 描述
Any

The resulting Any value.

源代码位于: jianmu/engine/state.py
def get_runtime_metadata(self, key: str, default: Any = None) -> Any:
    """Return one runner-internal metadata value.

    Args:
        key: The `key` value.
        default: Default value to use when no explicit value is available.

    Returns:
        The resulting `Any` value.
    """
    with self._lock:
        if key not in self._runtime_metadata:
            return copy.deepcopy(default)
        return copy.deepcopy(self._runtime_metadata[key])

clear_runtime_metadata

clear_runtime_metadata(key: str | None = None) -> None

Clear one runtime metadata key or the whole runtime metadata store.

参数:

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

The key value.

None
源代码位于: jianmu/engine/state.py
def clear_runtime_metadata(self, key: str | None = None) -> None:
    """Clear one runtime metadata key or the whole runtime metadata store.

    Args:
        key: The `key` value.
    """
    if not state_write_allowed():
        _guarded_write_blocked("clear_runtime_metadata")
        return
    with self._lock:
        if key is None:
            if not self._runtime_metadata:
                return
            self._runtime_metadata.clear()
        else:
            normalized_key = str(key)
            if normalized_key not in self._runtime_metadata:
                return
            self._runtime_metadata.pop(normalized_key, None)
        self._bump_revision()
    self.signal()

restore_runtime_metadata

restore_runtime_metadata(
    data: Dict[str, Any] | None,
) -> None

Replace runtime metadata from checkpoint payload.

参数:

名称 类型 描述 默认
data Dict[str, Any] | None

The data value.

必需
源代码位于: jianmu/engine/state.py
def restore_runtime_metadata(self, data: Dict[str, Any] | None) -> None:
    """Replace runtime metadata from checkpoint payload.

    Args:
        data: The `data` value.
    """
    with self._lock:
        next_metadata = copy.deepcopy(dict(data or {}))
        if self._runtime_metadata == next_metadata:
            return
        self._runtime_metadata = next_metadata
        self._bump_revision()

reset_step_fields

reset_step_fields()

Reset all scope='step' fields.

源代码位于: jianmu/engine/state.py
def reset_step_fields(self):
    """Reset all ``scope='step'`` fields."""
    self._reset_fields_by_scope("step")

reset_run_fields

reset_run_fields()

Reset all scope='run' fields.

源代码位于: jianmu/engine/state.py
def reset_run_fields(self):
    """Reset all ``scope='run'`` fields."""
    self._reset_fields_by_scope("run")

reset_call_fields

reset_call_fields()

Reset all scope='call' fields.

ReactiveRunner invokes this before one model/tool call cycle so temporary call-scoped state does not leak across turns.

源代码位于: jianmu/engine/state.py
def reset_call_fields(self):
    """Reset all ``scope='call'`` fields.

    ``ReactiveRunner`` invokes this before one model/tool call cycle so
    temporary call-scoped state does not leak across turns.
    """
    self._reset_fields_by_scope("call")

get_step_fields

get_step_fields() -> Dict[str, Any]

Return the current values of all scope='step' fields.

返回:

类型 描述
Dict[str, Any]

Mapping of step-scoped field names to their current values.

源代码位于: jianmu/engine/state.py
def get_step_fields(self) -> Dict[str, Any]:
    """Return the current values of all ``scope='step'`` fields.

    Returns:
        Mapping of step-scoped field names to their current values.
    """
    actions = {}
    with self._lock:
        if self._data is None:
            return actions

        data_dict = self._data.model_dump()
        for name, (_, _, scope) in self._ephemeral_fields.items():
            if scope == "step" and name in data_dict:
                actions[name] = data_dict[name]
    return actions

ReactiveRunner

ReactiveRunner

ReactiveRunner(
    root: Behaviour,
    state_manager: StateManager,
    ctx: Optional[RunContext] = None,
    *,
    setup_timeout: float = 15.0,
    max_pending_wakeups: int | None = None,
    hot_loop_warn_factor: float | None = None,
)

Event-driven runner that drives a Jianmu behaviour tree to completion.

The runner binds a root behaviour tree to a StateManager and optional RunContext. It injects runtime dependencies into Jianmu nodes, reacts to state changes, supports step mode for interactive apps, and can restore checkpoints for resumable execution.

The three most important execution entrypoints are:

  • tick_once(): perform a single synchronous tree tick
  • step(): one async tick plus step-scoped state handling
  • run(): drive the tree until it reaches SUCCESS or FAILURE

Minimal usage::

runner = ReactiveRunner(root=my_tree, state_manager=sm, ctx=ctx)
status = await runner.run({"task": "solve this problem"})

属性:

名称 类型 描述
root

Root behaviour bound to the runner.

state_manager

Shared runtime state store used by all injected nodes.

ctx

Active runtime dependency context.

tree

Wrapped py_trees behaviour tree.

Create a runner for a behaviour tree and state manager.

参数:

名称 类型 描述 默认
root Behaviour

Root behaviour node of the tree.

必需
state_manager StateManager

Durable state store shared across all nodes.

必需
ctx Optional[RunContext]

Optional runtime dependency context injected into Jianmu nodes.

None
setup_timeout float

Maximum setup time passed to py_trees before the tree is considered failed to initialize.

15.0
max_pending_wakeups int | None

Maximum buffered wakeups retained before coalescing repeated notifications.

None
hot_loop_warn_factor float | None

Warning threshold multiplier relative to max_fps for detecting pathological retry loops.

None
源代码位于: jianmu/engine/runtime.py
def __init__(
    self,
    root: py_trees.behaviour.Behaviour,
    state_manager: StateManager,
    ctx: Optional[RunContext] = None,
    *,
    setup_timeout: float = 15.0,
    max_pending_wakeups: int | None = None,
    hot_loop_warn_factor: float | None = None,
):
    """Create a runner for a behaviour tree and state manager.

    Args:
        root: Root behaviour node of the tree.
        state_manager: Durable state store shared across all nodes.
        ctx: Optional runtime dependency context injected into Jianmu nodes.
        setup_timeout: Maximum setup time passed to ``py_trees`` before the
            tree is considered failed to initialize.
        max_pending_wakeups: Maximum buffered wakeups retained before
            coalescing repeated notifications.
        hot_loop_warn_factor: Warning threshold multiplier relative to
            ``max_fps`` for detecting pathological retry loops.
    """
    from jianmu.config.loader import get_config

    runtime_config = get_config().runtime
    self.root = root
    self.state_manager = state_manager
    self.ctx = ctx or RunContext()
    self.tree = BehaviourTree(root)

    self.tick_signal = asyncio.Event()
    # Wakeups are advisory signals rather than transactional events. Allow
    # one in-flight wakeup plus one extra buffered edge while preventing
    # unbounded backlog growth under noisy producers.
    self._tick_lock = threading.Lock()
    self._pending_tick_count = 0
    resolved_max_pending_wakeups = (
        runtime_config.max_pending_wakeups
        if max_pending_wakeups is None
        else max_pending_wakeups
    )
    if resolved_max_pending_wakeups <= 0:
        raise ValueError("max_pending_wakeups must be > 0")
    self._max_pending_ticks = int(resolved_max_pending_wakeups)
    self._hot_loop_warn_factor = (
        runtime_config.hot_loop_warn_factor
        if hot_loop_warn_factor is None
        else hot_loop_warn_factor
    )
    if self._hot_loop_warn_factor <= 0:
        raise ValueError("hot_loop_warn_factor must be > 0")

    # Internal wakeups are only consumed during event-driven run mode.
    self.auto_driving = False
    self._mode: Literal["idle", "step", "run"] = "idle"

    payload = InjectPayload(
        context=self.ctx,
        state_manager=self.state_manager,
        wake_up=self._on_wake_signal,
    )

    for node in self.root.iterate():
        inject_runtime_deps(node, payload)

    self.tree.setup(timeout=setup_timeout)

    # Thread-safe context
    self._loop = None
    self._loop_thread = None
    self._loop_owner_token: _ExecutionToken | None = None
    self._last_suspended_event_key: tuple[str, str] | None = None
    self._live_approval_request_id: str | None = None
    self._live_approval_unsubscribe = None
    self._live_approval_owner_token: _ExecutionToken | None = None
    self._last_checkpoint_signatures: dict[str, str] = {}
    self._last_checkpoint_cursors: dict[str, tuple[int, str]] = {}
    self._execution_lock = threading.Lock()
    self._execution_generation = 0
    self._active_execution_token: _ExecutionToken | None = None

tick_once

tick_once() -> Status

Execute one synchronous tree tick without entering the run loop.

This is the lowest-level execution API. Prefer step() for interactive stepping and run() for normal agent execution.

返回:

类型 描述
Status

The Status of the root node after the tick.

源代码位于: jianmu/engine/runtime.py
def tick_once(self) -> Status:
    """Execute one synchronous tree tick without entering the run loop.

    This is the lowest-level execution API. Prefer ``step()`` for
    interactive stepping and ``run()`` for normal agent execution.

    Returns:
        The Status of the root node after the tick.
    """
    self.tree.tick()
    return self.root.status

step async

step(
    obs: Optional[Dict[str, Any]] = None,
    yield_to_async: bool = False,
) -> Dict[str, Any]

Run exactly one tree tick in step mode.

step() is designed for UIs or orchestrators that want explicit turn-by-turn control. It resets Ephemeral("step") fields, merges an optional observation payload, performs one tick, and returns the current step-scoped action fields.

参数:

名称 类型 描述 默认
obs Optional[Dict[str, Any]]

Optional observation data merged into state before ticking.

None
yield_to_async bool

Whether to yield once to the event loop after the tick.

False

返回:

类型 描述
Dict[str, Any]

Current step-scoped action payload from the state manager.

引发:

类型 描述
RuntimeError

If run() is currently active.

源代码位于: jianmu/engine/runtime.py
async def step(
    self,
    obs: Optional[Dict[str, Any]] = None,
    yield_to_async: bool = False,
) -> Dict[str, Any]:
    """Run exactly one tree tick in step mode.

    ``step()`` is designed for UIs or orchestrators that want explicit
    turn-by-turn control. It resets ``Ephemeral("step")`` fields, merges an
    optional observation payload, performs one tick, and returns the
    current step-scoped action fields.

    Args:
        obs: Optional observation data merged into state before ticking.
        yield_to_async: Whether to yield once to the event loop after the tick.

    Returns:
        Current step-scoped action payload from the state manager.

    Raises:
        RuntimeError: If run() is currently active.
    """
    if self._mode == "run":
        raise RuntimeError("Cannot step() while run() is active")

    self._mode = "step"
    self.auto_driving = False

    token = None
    try:
        try:
            state = self.state_manager.get()
            agent_id = getattr(state, "agent_id", None)
            session_id = getattr(state, "session_id", None)
            if agent_id or session_id:
                token = trace_set_context(agent_id=agent_id, session_id=session_id)
        except Exception:
            token = None

        self.state_manager.reset_step_fields()
        if obs:
            self.state_manager.update(obs)

        before_status = self.root.status.name if self.root.status is not None else None
        await self._notify_runtime_hook(
            "before_tick",
            event="tick.before",
            status=before_status,
            payload={"mode": "step"},
        )
        self.tick_once()
        await self._notify_runtime_hook(
            "after_tick",
            event="tick.after",
            status=self.root.status.name,
            payload={"mode": "step"},
        )

        if yield_to_async:
            await asyncio.sleep(0)

        return self.state_manager.get_step_fields()
    finally:
        self._mode = "idle"
        if token is not None:
            trace_reset_context(token)

run async

run(
    input_data: Optional[Dict[str, Any]] = None,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpointer: Any = None,
    checkpoint_interval: int | None = 1,
    thread_id: str = "default_thread",
    max_fps: float = 60.0,
) -> Status

Drive the behaviour tree until it finishes or hits a guardrail.

run() is the normal end-to-end execution mode. It can optionally reset the tree, reset or reuse state, seed initial input data, restore from a checkpoint backend, and then keep ticking until the root reaches a terminal status.

参数:

名称 类型 描述 默认
input_data Optional[Dict[str, Any]]

Optional initial input state applied before the run.

None
reset_tree bool

Whether to interrupt and reset tree execution first.

True
reset_data bool

Whether to fully reinitialize state before the run.

False
max_ticks int

Optional hard tick limit.

None
timeout_s float | None

Optional wall-clock timeout for the full run.

None
checkpointer Any

Optional checkpoint backend implementing Jianmu's checkpointer protocol.

None
checkpoint_interval int | None

Tick interval between checkpoint saves, or None to save only when checkpointable state/tree status changes.

1
thread_id str

Logical thread/session identifier for checkpoints.

'default_thread'
max_fps float

Maximum runner tick frequency.

60.0

返回:

类型 描述
Status

Final root node status.

引发:

类型 描述
RuntimeError

If step() is currently active.

源代码位于: jianmu/engine/runtime.py
async def run(
    self,
    input_data: Optional[Dict[str, Any]] = None,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpointer: Any = None,
    checkpoint_interval: int | None = 1,
    thread_id: str = "default_thread",
    max_fps: float = 60.0,
) -> Status:
    """Drive the behaviour tree until it finishes or hits a guardrail.

    ``run()`` is the normal end-to-end execution mode. It can optionally
    reset the tree, reset or reuse state, seed initial input data, restore
    from a checkpoint backend, and then keep ticking until the root reaches
    a terminal status.

    Args:
        input_data: Optional initial input state applied before the run.
        reset_tree: Whether to interrupt and reset tree execution first.
        reset_data: Whether to fully reinitialize state before the run.
        max_ticks: Optional hard tick limit.
        timeout_s: Optional wall-clock timeout for the full run.
        checkpointer: Optional checkpoint backend implementing Jianmu's
            checkpointer protocol.
        checkpoint_interval: Tick interval between checkpoint saves, or
            ``None`` to save only when checkpointable state/tree status
            changes.
        thread_id: Logical thread/session identifier for checkpoints.
        max_fps: Maximum runner tick frequency.

    Returns:
        Final root node status.

    Raises:
        RuntimeError: If step() is currently active.
    """
    self._enter_run_mode("run")
    execution_token = self._begin_execution()
    write_guard_token = push_state_write_guard(
        lambda: self._can_commit_execution(execution_token),
        info={
            "entrypoint": "run",
            "generation": execution_token.generation,
            "thread_id": thread_id,
            "runner_id": f"{id(self):x}",
        },
    )
    token = None
    try:
        try:
            state = self.state_manager.get()
            agent_id = getattr(state, "agent_id", None)
            session_id = getattr(state, "session_id", None)
            if agent_id or session_id:
                token = trace_set_context(agent_id=agent_id, session_id=session_id)
        except Exception:
            token = None

        # `_event_loop()` clears wake callbacks on exit to avoid long-lived
        # stale references, so each run must re-inject them.
        payload = InjectPayload(
            context=self.ctx,
            state_manager=self.state_manager,
            wake_up=lambda: self._on_wake_signal(execution_token),
        )
        for node in self.root.iterate():
            inject_runtime_deps(node, payload)

        if reset_tree:
            self.tree.interrupt()

        if reset_data:
            self.state_manager.initialize()
        else:
            self.state_manager.reset_run_fields()
            self.state_manager.update({"done": False}, signal=False)
        self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)

        self.tick_signal.clear()

        if input_data:
            self.state_manager.update(input_data)
        self.state_manager.set_runtime_metadata(InteractionKeys.THREAD_ID, thread_id)

        self._reset_run_counters()

        return await self._event_loop(
            max_ticks=max_ticks,
            timeout_s=timeout_s,
            checkpointer=checkpointer,
            checkpoint_interval=checkpoint_interval,
            thread_id=thread_id,
            max_fps=max_fps,
            execution_token=execution_token,
        )
    finally:
        self._mode = "idle"
        self._invalidate_execution_token(execution_token)
        pop_state_write_guard(write_guard_token)
        if token is not None:
            trace_reset_context(token)

run_until_suspend async

run_until_suspend(
    input_data: Optional[Dict[str, Any]] = None,
    *,
    resume_data: Optional[Dict[str, Any]] = None,
    resume_request_id: str | None = None,
    restore: str | RestorePolicy = RestorePolicy.NEVER,
    suspension_mode: str
    | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpointer: Any = None,
    checkpoint_interval: int | None = 1,
    thread_id: str = "default_thread",
    max_fps: float = 60.0,
) -> RunResult

Run until terminal completion or a publishable suspension point.

参数:

名称 类型 描述 默认
input_data Optional[Dict[str, Any]]

The input_data value.

None
resume_data Optional[Dict[str, Any]]

The resume_data value.

None
resume_request_id str | None

Active suspension id required when resume_data is provided.

None
restore str | RestorePolicy

The restore value.

NEVER
suspension_mode str | SuspensionMode

The suspension_mode value.

YIELD
reset_tree bool

The reset_tree value.

True
reset_data bool

The reset_data value.

False
max_ticks int

Collection of max tick values.

None
timeout_s float | None

Collection of timeout values.

None
checkpointer Any

The checkpointer value.

None
checkpoint_interval int | None

Tick interval between checkpoint saves, or None to save only when checkpointable state/tree status changes.

1
thread_id str

Identifier for thread.

'default_thread'
max_fps float

Collection of max fp values.

60.0

返回:

类型 描述
RunResult

The resulting RunResult value.

引发:

类型 描述
RuntimeError

If the operation cannot be completed at runtime.

ValueError

If input validation fails.

源代码位于: jianmu/engine/runtime.py
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
1567
1568
1569
1570
1571
1572
1573
1574
1575
1576
1577
1578
1579
1580
1581
1582
1583
1584
1585
1586
1587
1588
1589
1590
1591
1592
1593
1594
1595
1596
1597
1598
1599
1600
1601
1602
1603
1604
1605
1606
1607
1608
1609
1610
1611
1612
1613
1614
1615
1616
1617
1618
1619
1620
1621
1622
1623
1624
1625
1626
1627
1628
1629
1630
1631
1632
1633
1634
1635
1636
1637
1638
1639
1640
1641
1642
1643
1644
1645
1646
1647
1648
1649
1650
1651
1652
1653
1654
1655
1656
1657
1658
1659
1660
1661
1662
1663
1664
1665
1666
1667
1668
1669
1670
1671
1672
1673
1674
1675
1676
1677
1678
1679
1680
1681
1682
1683
1684
1685
1686
1687
1688
1689
1690
1691
1692
1693
1694
1695
1696
1697
1698
1699
1700
1701
1702
1703
1704
1705
1706
1707
1708
1709
1710
1711
1712
1713
1714
1715
1716
1717
1718
1719
1720
1721
1722
1723
1724
1725
1726
1727
1728
1729
1730
1731
1732
1733
1734
1735
1736
1737
1738
1739
1740
1741
1742
1743
1744
1745
1746
1747
1748
1749
1750
1751
1752
1753
1754
1755
1756
1757
1758
1759
1760
1761
1762
1763
1764
1765
1766
1767
1768
1769
1770
1771
1772
1773
1774
1775
1776
1777
1778
1779
1780
1781
1782
1783
1784
1785
1786
1787
1788
1789
1790
1791
1792
1793
1794
1795
1796
1797
1798
1799
1800
1801
1802
1803
1804
1805
1806
1807
1808
1809
1810
1811
1812
1813
1814
1815
1816
1817
1818
1819
1820
1821
1822
1823
1824
1825
1826
1827
1828
1829
1830
1831
1832
1833
1834
1835
1836
1837
1838
1839
1840
1841
1842
1843
1844
1845
1846
1847
1848
1849
1850
1851
1852
1853
1854
1855
1856
1857
1858
1859
1860
1861
1862
1863
1864
1865
1866
1867
1868
1869
1870
1871
1872
1873
1874
1875
1876
1877
1878
1879
1880
1881
1882
async def run_until_suspend(
    self,
    input_data: Optional[Dict[str, Any]] = None,
    *,
    resume_data: Optional[Dict[str, Any]] = None,
    resume_request_id: str | None = None,
    restore: str | RestorePolicy = RestorePolicy.NEVER,
    suspension_mode: str | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpointer: Any = None,
    checkpoint_interval: int | None = 1,
    thread_id: str = "default_thread",
    max_fps: float = 60.0,
) -> RunResult:
    """Run until terminal completion or a publishable suspension point.

    Args:
        input_data: The `input_data` value.
        resume_data: The `resume_data` value.
        resume_request_id: Active suspension id required when resume_data is provided.
        restore: The `restore` value.
        suspension_mode: The `suspension_mode` value.
        reset_tree: The `reset_tree` value.
        reset_data: The `reset_data` value.
        max_ticks: Collection of max tick values.
        timeout_s: Collection of timeout  values.
        checkpointer: The `checkpointer` value.
        checkpoint_interval: Tick interval between checkpoint saves, or
            ``None`` to save only when checkpointable state/tree status
            changes.
        thread_id: Identifier for thread.
        max_fps: Collection of max fp values.

    Returns:
        The resulting `RunResult` value.

    Raises:
        RuntimeError: If the operation cannot be completed at runtime.
        ValueError: If input validation fails.
    """
    if resume_data is not None and resume_request_id is None:
        raise ValueError("resume_request_id is required when resume_data is provided.")
    restore_policy = restore if isinstance(restore, RestorePolicy) else RestorePolicy(str(restore))
    mode = suspension_mode if isinstance(suspension_mode, SuspensionMode) else SuspensionMode(str(suspension_mode))
    self._enter_run_mode("run_until_suspend")
    execution_token = self._begin_execution()
    write_guard_token = push_state_write_guard(
        lambda: self._can_commit_execution(execution_token),
        info={
            "entrypoint": "run_until_suspend",
            "generation": execution_token.generation,
            "thread_id": thread_id,
            "runner_id": f"{id(self):x}",
        },
    )
    token = None
    try:
        try:
            state = self.state_manager.get()
            agent_id = getattr(state, "agent_id", None)
            session_id = getattr(state, "session_id", None)
            if agent_id or session_id:
                token = trace_set_context(agent_id=agent_id, session_id=session_id)
        except Exception:
            token = None

        payload = InjectPayload(
            context=self.ctx,
            state_manager=self.state_manager,
            wake_up=lambda: self._on_wake_signal(execution_token),
        )
        for node in self.root.iterate():
            inject_runtime_deps(node, payload)

        runtime = getattr(self.ctx, "runtime", None) if self.ctx is not None else None
        if runtime is not None and getattr(runtime, "thread_id", None) != thread_id:
            try:
                runtime.thread_id = thread_id
            except Exception:
                pass

        if reset_tree:
            self.tree.interrupt()

        if reset_data:
            self.state_manager.initialize()
        else:
            self.state_manager.reset_run_fields()
            self.state_manager.update({"done": False}, signal=False)
        self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)

        self.tick_signal.clear()

        interaction = InteractionController(self.state_manager)
        restored = False
        if restore_policy != RestorePolicy.NEVER:
            restored = self._try_restore_from_checkpoint(
                checkpointer,
                thread_id,
                require_active_path_match=restore_policy == RestorePolicy.REQUIRED,
            )
            if restore_policy == RestorePolicy.REQUIRED and not restored:
                if resume_request_id is not None:
                    result = ResumeSuspensionResult(
                        status=ResumeStatus.INVALID_REQUEST,
                        request_id=resume_request_id,
                        reason="checkpoint_not_found",
                        message=f"Checkpoint required for thread_id={thread_id!r} but none was found.",
                        payload={"thread_id": thread_id},
                        retryable=False,
                    )
                    raise_for_resume_result(result)
                raise ValueError(f"Checkpoint required for thread_id={thread_id!r} but none was found.")
            if restored:
                termination = TerminationRecord.from_dict(
                    self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
                )
                if termination is not None and termination.reason == "checkpoint_incompatible":
                    return await self._return_run_result(
                        RunResult(
                            outcome="failure",
                            status=Status.FAILURE.name,
                            termination=termination,
                        ),
                        event="run.failed",
                        payload={"thread_id": thread_id},
                    )

        if restore_policy == RestorePolicy.NEVER and input_data:
            self.state_manager.update(input_data)
        elif input_data and not restored:
            self.state_manager.update(input_data)
        self.state_manager.set_runtime_metadata(InteractionKeys.THREAD_ID, thread_id)

        if restored:
            active = interaction.get_active_suspension()
            if active is not None and resume_data is None:
                raise ValueError(
                    "Resume data is required when restoring an active suspended workflow."
                )
        resumed_suspension = interaction.get_active_suspension()
        if resume_request_id is not None:
            result = build_resume_data(
                suspension=resumed_suspension,
                request_id=resume_request_id,
                payload=resume_data,
            )
            raise_for_resume_result(result)
            resume_data = result.payload
        if resume_data is not None:
            if resumed_suspension is None:
                resumed_suspension = interaction.get_active_suspension()
            resume_kind = "payload"
            if (
                resumed_suspension is not None
                and resumed_suspension.category == SuspensionCategory.REQUIRE_APPROVAL
            ):
                approval_result = resolve_approval_resume(
                    manager=getattr(self.ctx, "approval_manager", None) if self.ctx is not None else None,
                    request_id=resumed_suspension.request_id,
                    category=resumed_suspension.category.value,
                    payload=resume_data,
                )
                raise_for_resume_result(self._approval_resume_to_runtime_result(approval_result))
                interaction.write_resolved_approval(
                    resumed_suspension.request_id,
                    dict(approval_result.decision or {}),
                )
                resume_kind = "approval"
            else:
                interaction.write_resume_payload(resume_data, signal=False)
            interaction.deactivate_suspension(signal=False)
            self._last_suspended_event_key = None
            resumed_payload = {
                "thread_id": thread_id,
                "lifecycle_category": "interaction",
                "resume_kind": resume_kind,
                "payload": copy.deepcopy(dict(resumed_suspension.payload or {}))
                if resumed_suspension is not None
                else {},
            }
            if resumed_suspension is not None:
                resumed_payload.update(
                    {
                        "request_id": resumed_suspension.request_id,
                        "reason": resumed_suspension.reason.value,
                        "category": resumed_suspension.category.value,
                    }
                )
            self._emit_runtime_event(
                "execution.resumed",
                payload=resumed_payload,
            )
            await self._notify_runtime_hook(
                "on_resume",
                event="execution.resumed",
                status=self.root.status.name,
                payload=resumed_payload,
            )

        self._reset_run_counters()

        if max_fps <= 0:
            raise ValueError(f"max_fps must be > 0, got {max_fps}")
        if checkpoint_interval is not None and checkpoint_interval <= 0:
            raise ValueError(f"checkpoint_interval must be > 0, got {checkpoint_interval}")
        if timeout_s is not None and timeout_s <= 0:
            raise ValueError(f"timeout_s must be > 0, got {timeout_s}")
        min_tick_interval = 1.0 / max_fps
        self._bind_callbacks(execution_token)
        self._bind_execution_loop(execution_token)
        self.auto_driving = True
        self._signal_tick(execution_token)
        total_tick_count = 0
        run_started_at = time.monotonic()
        run_deadline = run_started_at + timeout_s if timeout_s is not None else None
        try:
            while True:
                if timeout_s is not None and (time.monotonic() - run_started_at) >= timeout_s:
                    status = self.root.status
                    if status in (Status.SUCCESS, Status.FAILURE):
                        if status == Status.SUCCESS:
                            self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)
                            return await self._return_run_result(
                                RunResult(outcome="success", status=status.name),
                                event="run.completed",
                                payload={"thread_id": thread_id},
                            )
                        termination = TerminationRecord.from_dict(
                            self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
                        )
                        return await self._finalize_failure_result(
                            status=status,
                            termination=termination,
                            thread_id=thread_id,
                        )
                    status = Status.FAILURE
                    termination = TerminationRecord(
                        reason="timeout_exceeded",
                        message=f"Reached timeout_s={timeout_s} before the tree finished.",
                        payload={"timeout_s": timeout_s, "elapsed_s": time.monotonic() - run_started_at},
                    )
                    self.state_manager.set_runtime_metadata(InteractionKeys.TERMINATION, termination.to_dict())
                    self._emit_runtime_event(
                        "execution.timeout_exceeded",
                        payload={
                            "thread_id": thread_id,
                            "timeout_s": timeout_s,
                            "elapsed_s": time.monotonic() - run_started_at,
                        },
                    )
                    return await self._return_run_result(
                        RunResult(outcome="failure", status=status.name, termination=termination),
                        event="run.failed",
                        payload={"thread_id": thread_id},
                    )
                if max_ticks is not None and total_tick_count >= max_ticks:
                    status = self.root.status
                    if status in (Status.SUCCESS, Status.FAILURE):
                        if status == Status.SUCCESS:
                            self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)
                            return await self._return_run_result(
                                RunResult(outcome="success", status=status.name),
                                event="run.completed",
                                payload={"thread_id": thread_id},
                            )
                        termination = TerminationRecord.from_dict(
                            self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
                        )
                        return await self._finalize_failure_result(
                            status=status,
                            termination=termination,
                            thread_id=thread_id,
                        )
                    status = Status.FAILURE
                    termination = TerminationRecord(
                        reason="max_ticks_exceeded",
                        message=f"Reached max_ticks={max_ticks} before the tree finished.",
                        payload={"max_ticks": max_ticks, "tick_count": total_tick_count},
                    )
                    self.state_manager.set_runtime_metadata(InteractionKeys.TERMINATION, termination.to_dict())
                    self._emit_runtime_event(
                        "execution.max_ticks_exceeded",
                        payload={
                            "thread_id": thread_id,
                            "max_ticks": max_ticks,
                            "tick_count": total_tick_count,
                        },
                    )
                    return await self._return_run_result(
                        RunResult(outcome="failure", status=status.name, termination=termination),
                        event="run.failed",
                        payload={"thread_id": thread_id},
                    )

                try:
                    await self._wait_for_tick(deadline=run_deadline)
                except asyncio.TimeoutError:
                    status = self.root.status
                    if status in (Status.SUCCESS, Status.FAILURE):
                        if status == Status.SUCCESS:
                            self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)
                            return await self._return_run_result(
                                RunResult(outcome="success", status=status.name),
                                event="run.completed",
                                payload={"thread_id": thread_id},
                            )
                        termination = TerminationRecord.from_dict(
                            self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
                        )
                        return await self._finalize_failure_result(
                            status=status,
                            termination=termination,
                            thread_id=thread_id,
                        )
                    status = Status.FAILURE
                    termination = TerminationRecord(
                        reason="timeout_exceeded",
                        message=f"Reached timeout_s={timeout_s} before the tree finished.",
                        payload={"timeout_s": timeout_s, "elapsed_s": time.monotonic() - run_started_at},
                    )
                    self.state_manager.set_runtime_metadata(InteractionKeys.TERMINATION, termination.to_dict())
                    self._emit_runtime_event(
                        "execution.timeout_exceeded",
                        payload={
                            "thread_id": thread_id,
                            "timeout_s": timeout_s,
                            "elapsed_s": time.monotonic() - run_started_at,
                        },
                    )
                    return await self._return_run_result(
                        RunResult(outcome="failure", status=status.name, termination=termination),
                        event="run.failed",
                        payload={"thread_id": thread_id},
                    )
                tick_start_time = time.monotonic()
                self._sync_live_approval_subscription(
                    interaction=interaction,
                    thread_id=thread_id,
                    execution_token=execution_token,
                )
                self._resume_live_approval_if_ready(
                    interaction=interaction,
                    thread_id=thread_id,
                    execution_token=execution_token,
                )
                before_status = self.root.status.name if self.root.status is not None else None
                await self._notify_runtime_hook(
                    "before_tick",
                    event="tick.before",
                    status=before_status,
                    payload={
                        "mode": "run_until_suspend",
                        "thread_id": thread_id,
                        "tick_count": total_tick_count,
                    },
                )
                self.tree.tick()
                total_tick_count += 1
                status = self.root.status
                await self._notify_runtime_hook(
                    "after_tick",
                    event="tick.after",
                    status=status.name,
                    payload={
                        "mode": "run_until_suspend",
                        "thread_id": thread_id,
                        "tick_count": total_tick_count,
                    },
                )
                tick_elapsed = time.monotonic() - tick_start_time
                if tick_elapsed < min_tick_interval:
                    await asyncio.sleep(min_tick_interval - tick_elapsed)
                else:
                    await asyncio.sleep(0)

                suspension = interaction.get_active_suspension()
                if suspension is not None and status == Status.RUNNING:
                    suspended_key = (thread_id, suspension.request_id)
                    if self._last_suspended_event_key != suspended_key:
                        self._emit_runtime_event(
                            "execution.suspended",
                            payload={
                                "lifecycle_category": "interaction",
                                "reason": suspension.reason.value,
                                "category": suspension.category.value,
                                "request_id": suspension.request_id,
                                "message": suspension.message,
                                "thread_id": thread_id,
                                "payload": copy.deepcopy(dict(suspension.payload or {})),
                            },
                        )
                        self._last_suspended_event_key = suspended_key
                    if mode == SuspensionMode.YIELD:
                        self._save_checkpoint_now(
                            checkpointer,
                            thread_id,
                            total_tick_count,
                        )
                        suspended_result = RunResult(
                            outcome="suspended",
                            status=status.name,
                            suspension=suspension,
                        )
                        await self._notify_runtime_hook(
                            "on_suspend",
                            event="execution.suspended",
                            status=status.name,
                            payload={
                                "thread_id": thread_id,
                                "reason": suspension.reason.value,
                                "category": suspension.category.value,
                                "request_id": suspension.request_id,
                                "message": suspension.message,
                                "payload": copy.deepcopy(dict(suspension.payload or {})),
                            },
                        )
                        return suspended_result

                self._maybe_save_checkpoint(
                    checkpointer=checkpointer,
                    checkpoint_interval=checkpoint_interval,
                    thread_id=thread_id,
                    total_tick_count=total_tick_count,
                )

                if status == Status.SUCCESS:
                    self.state_manager.clear_runtime_metadata(InteractionKeys.TERMINATION)
                    self._flush_acknowledged_approval_results(
                        checkpointer=checkpointer,
                        thread_id=thread_id,
                        step=total_tick_count,
                    )
                    return await self._return_run_result(
                        RunResult(outcome="success", status=status.name),
                        event="run.completed",
                        payload={"thread_id": thread_id},
                    )
                if status == Status.FAILURE:
                    termination = TerminationRecord.from_dict(
                        self.state_manager.get_runtime_metadata(InteractionKeys.TERMINATION)
                    )
                    self._flush_acknowledged_approval_results(
                        checkpointer=checkpointer,
                        thread_id=thread_id,
                        step=total_tick_count,
                    )
                    return await self._finalize_failure_result(
                        status=status,
                        termination=termination,
                        thread_id=thread_id,
                    )
        finally:
            self._cleanup_execution_resources(execution_token)
            if sys.exc_info()[0] is None:
                self._flush_acknowledged_approval_results(
                    checkpointer=checkpointer,
                    thread_id=thread_id,
                    step=total_tick_count,
                )
            else:
                self._flush_acknowledged_approval_results_best_effort(
                    checkpointer=checkpointer,
                    thread_id=thread_id,
                    step=total_tick_count,
                )
    finally:
        self._mode = "idle"
        self._invalidate_execution_token(execution_token)
        pop_state_write_guard(write_guard_token)
        if token is not None:
            trace_reset_context(token)

resume_and_continue async

resume_and_continue(
    *,
    thread_id: str,
    request_id: str,
    payload: dict[str, Any],
    checkpointer: Any = None,
    suspension_mode: str
    | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpoint_interval: int | None = 1,
    max_fps: float = 60.0,
    **kwargs: Any,
) -> RunResult

Resume a checkpointed suspension and continue execution.

This convenience entrypoint resumes a checkpointed suspension and then continues execution through run_until_suspend(). User-input, external-execution, and approval suspensions all flow through the same strict request_id validation path, while approval resumes bridge through the shared RunContext.approval_manager before tree re-entry.

参数:

名称 类型 描述 默认
thread_id str

Logical checkpoint thread identifier.

必需
request_id str

Active suspension request identifier to resume.

必需
payload dict[str, Any]

Host-provided resume payload.

必需
checkpointer Any

Optional checkpoint backend implementing Jianmu's protocol.

None
suspension_mode str | SuspensionMode

Host-facing handling mode for any subsequent suspension.

YIELD
reset_tree bool

Whether to interrupt and reset tree execution first.

True
reset_data bool

Whether to fully reinitialize state before the run.

False
max_ticks int

Optional hard tick limit.

None
timeout_s float | None

Optional wall-clock timeout for the resumed run.

None
checkpoint_interval int | None

Tick interval between checkpoint saves, or None to save only when checkpointable state/tree status changes.

1
max_fps float

Maximum runner tick frequency.

60.0
**kwargs Any

Additional keyword arguments forwarded to run_until_suspend.

{}

返回:

类型 描述
RunResult

The resumed run result.

源代码位于: jianmu/engine/runtime.py
async def resume_and_continue(
    self,
    *,
    thread_id: str,
    request_id: str,
    payload: dict[str, Any],
    checkpointer: Any = None,
    suspension_mode: str | SuspensionMode = SuspensionMode.YIELD,
    reset_tree: bool = True,
    reset_data: bool = False,
    max_ticks: int = None,
    timeout_s: float | None = None,
    checkpoint_interval: int | None = 1,
    max_fps: float = 60.0,
    **kwargs: Any,
) -> RunResult:
    """Resume a checkpointed suspension and continue execution.

    This convenience entrypoint resumes a checkpointed suspension and then
    continues execution through ``run_until_suspend()``. User-input,
    external-execution, and approval suspensions all flow through the same
    strict ``request_id`` validation path, while approval resumes bridge
    through the shared ``RunContext.approval_manager`` before tree re-entry.

    Args:
        thread_id: Logical checkpoint thread identifier.
        request_id: Active suspension request identifier to resume.
        payload: Host-provided resume payload.
        checkpointer: Optional checkpoint backend implementing Jianmu's protocol.
        suspension_mode: Host-facing handling mode for any subsequent suspension.
        reset_tree: Whether to interrupt and reset tree execution first.
        reset_data: Whether to fully reinitialize state before the run.
        max_ticks: Optional hard tick limit.
        timeout_s: Optional wall-clock timeout for the resumed run.
        checkpoint_interval: Tick interval between checkpoint saves, or
            ``None`` to save only when checkpointable state/tree status
            changes.
        max_fps: Maximum runner tick frequency.
        **kwargs: Additional keyword arguments forwarded to ``run_until_suspend``.

    Returns:
        The resumed run result.
    """
    return await self.run_until_suspend(
        resume_data=payload,
        resume_request_id=request_id,
        restore=RestorePolicy.REQUIRED,
        suspension_mode=suspension_mode,
        reset_tree=reset_tree,
        reset_data=reset_data,
        max_ticks=max_ticks,
        timeout_s=timeout_s,
        checkpointer=checkpointer,
        checkpoint_interval=checkpoint_interval,
        thread_id=thread_id,
        max_fps=max_fps,
        **kwargs,
    )

reset

reset(reset_data: bool = True)

Interrupt the tree and reset runner-local execution state.

参数:

名称 类型 描述 默认
reset_data bool

Whether to reinitialize the backing state manager.

True
源代码位于: jianmu/engine/runtime.py
def reset(self, reset_data: bool = True):
    """Interrupt the tree and reset runner-local execution state.

    Args:
        reset_data: Whether to reinitialize the backing state manager.
    """
    active_token = self._active_execution_token
    self._invalidate_execution_token(active_token)
    self._unsubscribe_state_wake_callback(active_token)
    self._clear_live_approval_subscription(force=True)
    self._clear_bound_execution_loop(active_token)
    self.tree.interrupt()
    if reset_data:
        self.state_manager.initialize()
    self.tick_signal.clear()
    self.auto_driving = False
    self._mode = "idle"
    self._last_suspended_event_key = None
    with self._tick_lock:
        self._pending_tick_count = 0
    clear_payload = InjectPayload(context=self.ctx, state_manager=self.state_manager, wake_up=None)
    for node in self.root.iterate():
        inject_runtime_deps(node, clear_payload)

RunContext

RunContext dataclass

RunContext(
    model_client: Any = None,
    approval_manager: Optional["ApprovalManager"] = None,
    sandbox: Optional[Any] = None,
    constraints: Optional[Constraints] = None,
    runtime: Optional[Any] = None,
    prompt_runtime: Optional[PromptRuntimeContext] = None,
    runtime_event_bus: Optional["RuntimeEventBus"] = None,
    effect_manager: Optional["EffectCommitManager"] = None,
    hooks: Optional["HookManager"] = None,
)

Runtime-only dependency container injected into executing nodes.

RunContext carries resources that should not live in serializable agent state, such as model clients, approval managers, or sandbox handles. ReactiveRunner injects the same context into every Jianmu-aware node before execution.

Typical usage keeps durable data in StateManager and places shared service objects here::

ctx = RunContext(
    model_client=my_model_client,
    constraints=my_constraints,
    approval_manager=my_approval_manager,
)

属性:

名称 类型 描述
model_client Any

Default model client used by nodes that resolve model execution capability from runtime context.

approval_manager Optional['ApprovalManager']

Optional approval workflow manager used by guarded tool execution.

sandbox Optional[Any]

Optional sandbox/runtime handle shared by execution nodes.

constraints Optional[Constraints]

Optional global constraints applied by guard/execution components.

runtime Optional[Any]

Optional opaque runtime extension slot for app-specific data.

prompt_runtime Optional[PromptRuntimeContext]

Optional prompt-construction overrides for the current run, such as bootstrap files or skill-aware prompt assembly.

runtime_event_bus Optional['RuntimeEventBus']

Optional runtime-semantic event bus consumed by nodes, guards, and runtimes to publish execution facts.

effect_manager Optional['EffectCommitManager']

Optional governed side-effect commit manager used by nodes/adapters that opt into prepare/commit/receipt execution.

hooks Optional['HookManager']

Optional hook manager for execution-scoped extension points.