State cn
Jianmu 的状态管理围绕一个核心理念构建:Agent 的整个运行时状态应该是一个带类型的、可验证的 Pydantic 模型实例,而不仅仅是无结构的字典。StateManager 作为这一理念的执行者,提供了线程安全的读写、字段级 Reducer 合并、自动复位 Ephemeral 字段,以及变更通知——这些机制共同支撑了行为树执行与 UI 交互之间的数据契约。
架构全景:StateManager 的三层写入与生命周期绑定¶
StateManager 不是一个简单的键值存储;它是一台精密的状态机,位于 Pydantic 验证层和 ReactiveRunner 事件循环之间。从架构上看,它提供三层写入 API,每一层对应不同的并发安全保证和语义:
graph TD
subgraph "用户定义层"
SCHEMA["Pydantic BaseModel<br/>• 普通字段: int, str, dict<br/>• Reducer 字段: Annotated[T, fn]<br/>• Ephemeral 字段: Annotated[T, Ephemeral(scope)]"]
end
subgraph "StateManager 核心"
PARSE["_parse_schema()<br/>提取 reducer + ephemeral 元数据"]
LOCK["threading.Lock<br/>全局读写锁"]
FW["per-field write locks<br/>序列化 read-modify-write"]
VALIDATE["Pydantic 验证<br/>每次写入后重建模型实例"]
NOTIFY["_notify_listeners()<br/>唤醒 ReactiveRunner"]
end
subgraph "三层写入 API"
UPDATE["update(dict)<br/>通用字段合并 + Reducer"]
MERGE["merge_dict(field, dict)<br/>原子字典合并"]
TRANSFORM["transform_field(key, fn)<br/>通用 read-modify-write"]
end
subgraph "生命周期绑定"
RUN["ReactiveRunner.run()<br/>→ reset_run_fields()"]
STEP["ReactiveRunner.step()<br/>→ reset_step_fields()"]
CALL["LLM Node update_async()<br/>→ reset_call_fields()"]
end
subgraph "消费者"
NODES["Jianmu 节点<br/>read_state / write_state / write_port"]
INTERACTION["InteractionController<br/>namespace='interaction'"]
MESSAGE_STORE["StateMessageStore<br/>transform_field 原子追加"]
end
SCHEMA --> PARSE --> LOCK
LOCK --> UPDATE & MERGE & TRANSFORM
UPDATE & MERGE & TRANSFORM --> VALIDATE --> NOTIFY
RUN & STEP & CALL -.->|Ephemeral 复位| LOCK
NOTIFY --> NODES & INTERACTION & MESSAGE_STORE
StateManager 初始化时接收一个 Pydantic 模型类型,然后通过 _parse_schema() 扫描所有字段的类型注解,分别提取 Reducer 函数和 Ephemeral 标记。这个解析过程发生在构造阶段,后续每次写入都会利用这些预提取的元数据来决定合并策略。
Pydantic Schema:定义 Agent 的类型化状态表面¶
Jianmu 不强制使用内置状态模型。用户通过继承 pydantic.BaseModel 定义自己的状态结构,然后将其传递给 StateManager 或 Agent 外观。这带来了三项关键收益:编译时类型安全(通过 IDE 自动补全)、运行时验证(Pydantic 在每次写入时重建模型实例)、默认值语义(未初始化的状态可以从 schema 默认值安全构造)。
一个典型的 ReAct Agent 状态 Schema 如下:
from typing import Annotated, List, Optional
from pydantic import BaseModel, Field
from jianmu.engine.state import Ephemeral
from jianmu.message import Message
from jianmu.tool.types import ToolCall
import operator
class ReActState(BaseModel):
messages: Annotated[List[Message], operator.add] = Field(default_factory=list)
done: Annotated[bool, Ephemeral(scope="run")] = False
final_answer: Annotated[Optional[str], Ephemeral(scope="run")] = None
actions: Annotated[List[ToolCall], Ephemeral(scope="step")] = Field(default_factory=list)
rounds: Annotated[int, Ephemeral(scope="run")] = 0
streaming_output: Annotated[str, Ephemeral(scope="call")] = ""
这个 Schema 同时展示了三种字段类别:messages 使用 operator.add 作为 Reducer——每次 update({"messages": [new_msg]}) 时,新消息会被追加到现有列表而非覆盖;done、final_answer、rounds 标记为 Ephemeral("run")——每次 run() 启动时自动复位;actions 标记为 Ephemeral("step")——每次 step() 执行前复位;streaming_output 标记为 Ephemeral("call")——每次 LLM 调用前清空。
Reducer 合并:声明式字段聚合策略¶
Reducer 是 Jianmu 状态管理的核心创新之一。它是一个绑定到特定字段的可调用对象,签名为 (old_value: T, new_value: T) -> T。当通过 update() 写入一个注册了 Reducer 的字段时,StateManager 不会简单地覆盖旧值,而是将旧值和新值同时传递给 Reducer,并将返回值作为最终写入值。
解析机制:_parse_schema() 使用 typing.get_type_hints(include_extras=True) 获取每个字段的完整注解。当发现注解来自 Annotated 时,它会遍历 Annotated 的额外参数:遇到 Ephemeral 实例则注册为临时字段,遇到其他可调用对象则注册为 Reducer。
工作流如下:
sequenceDiagram
participant Caller as 调用方
participant SM as StateManager
participant Reducer as Reducer 函数
participant Pydantic as Pydantic 验证
Caller->>SM: update({"messages": [msg_new]})
SM->>SM: 获取当前 model_dump()
SM->>SM: 检查 "messages" in reducers?
alt 有 Reducer
SM->>Reducer: reducer(old_val, new_val)
Reducer-->>SM: merged_value
else 无 Reducer
SM->>SM: final_val = new_val (覆盖)
end
SM->>SM: merged_data = {**current, field: final_val}
SM->>Pydantic: schema(**merged_data)
Pydantic-->>SM: 验证后的模型实例
SM->>Caller: 写入成功
SM->>SM: _notify_listeners()
最常用的 Reducer 是 operator.add,它对列表类型实现追加语义。测试清楚地验证了这一行为:连续两次 update({"history": ["Msg1"]}) 和 update({"history": ["Msg2"]}) 后,history 字段的值是 ["Init", "Msg1", "Msg2"] 而非仅 ["Msg2"]。
Reducer 的一个关键细节是首写容错:当 StateManager 尚未初始化(_data is None)且 Reducer 需要 old_val 时,系统会尝试从 schema 的字段默认值中获取 old_val,使得未调用 initialize() 时的首次 update() 也能正确工作。
字典合并与原子变换:merge_dict 和 transform_field¶
除了通用的 update(),StateManager 提供了两个更专用的写入 API,解决不同的并发场景:
| API | 语义 | 适用场景 | 并发安全 |
|---|---|---|---|
update(dict) |
字段级合并,支持 Reducer | 通用状态写入 | 全局锁 |
merge_dict(field, updates) |
原子字典合并,last-write-wins | 多个节点并发写入同一个共享 dict 字段 | 全局锁 |
transform_field(key, fn) |
原子 read-modify-write | 消息追加、计数器递增 | per-field 写锁 |
merge_dict 专为多节点并发输出设计。当多个并行节点把各自的输出写入同一个共享字典字段的不同键时,merge_dict 执行 {**current, **updates} 合并,确保并发的非重叠键写入不会互相覆盖。
并发安全性在测试中得到了验证:两个线程分别写入 {"a": 1} 和 {"b": 2},最终字典同时包含两个键。
transform_field 是最安全的写入方式。它获取字段级别的写锁(_field_write_locks),在锁保护下执行 fn(current_value),然后将返回值写回并验证。StateMessageStore 使用它来实现消息的原子追加:
self._state_manager.transform_field(
self._key,
lambda current: list(current or []) + [msg],
namespace=self._namespace,
)
这种设计消除了经典的"读取-修改-写入"竞态条件,在行为树并行节点环境中尤为重要。
Ephemeral 字段:三作用域自动复位¶
Ephemeral 是一个标记类(非 Pydantic 字段类型),它通过 typing.Annotated 附着在字段上,告诉 StateManager 该字段应该在特定生命周期节点自动复位到默认值。
stateDiagram-v2
[*] --> Idle
state "run() 入口" as RunEntry
state "step() 入口" as StepEntry
state "LLM Node" as LLMCall
Idle --> RunEntry: 用户调用 agent.run()
RunEntry --> Running: reset_run_fields()<br/>done=False, final_answer="", rounds=0
Idle --> StepEntry: 用户调用 agent.step()
StepEntry --> Ticking: reset_step_fields()<br/>actions=[], speed=0.0
Running --> LLMCall: LLM 节点开始调用
LLMCall --> Running: reset_call_fields()<br/>streaming_output=""
Running --> [*]: 运行结束
三个作用域的精确定义与触发时机:
| 作用域 | 触发时机 | 调用位置 | 典型字段 |
|---|---|---|---|
"run" |
ReactiveRunner.run() 和 run_until_suspend() 入口 |
runtime.py L737, L842 |
done, final_answer, rounds, tool_effects |
"step" |
ReactiveRunner.step() 入口,每次单步 tick 前 |
runtime.py L654 |
actions, speed, fire |
"call" |
AgentLLMNode.update_async() 和 SimpleLLMNode.update_async() 开始时 |
llm.py L152, L373 |
streaming_output |
复位机制:_reset_fields_by_scope(target_scope) 遍历所有注册的 Ephemeral 字段,跳过不匹配的作用域,然后对匹配字段执行:如果字段有 default_factory(如 Field(default_factory=list)),调用工厂函数获取新默认值;否则使用静态 default 值。复位还额外处理了 namespaced 变体——任何以 .{field_name} 结尾的键也会被同步复位。
step 作用域的字段还有特殊的 get_step_fields() 方法,供 ReactiveRunner.step() 返回当前 step 作用域字段的快照给 UI 层。这使得 TUI Chat 等交互界面可以在每次 step 后获取"本轮动作"而无需遍历完整状态。
Namespace:运行时扩展字段的安全隔离¶
当 Schema 配置了 model_config = ConfigDict(extra="allow") 后,StateManager 的 namespace 机制允许在 Pydantic 模型的主字段之外存储前缀化键值对。namespace 写入会被序列化为 "{namespace}.{key}" 格式并存储在 model_extra 中。
Namespace 有两个硬性约束:第一,namespace 写入不能与声明的 schema 字段同名——这会绕过 Reducer 和验证预期;第二,Schema 必须设置 extra='allow',否则严格模式会拒绝未定义字段。
InteractionController 是 namespace 的主要消费者。它使用 namespace="interaction" 存储挂起/恢复相关的运行时数据——suspension(当前挂起记录)、resume_payload(恢复载荷)、以及通过 runtime_metadata 管理的审批状态。这些数据不应该污染用户的业务状态字段,namespace 提供了自然的隔离边界。
clear(namespace) 方法支持批量删除某个 namespace 下的所有键,不影响其他 namespace 或主字段。
Listener 通知机制与 ReactiveRunner 联动¶
StateManager 维护一个回调列表,每次成功的写入(update、merge_dict、transform_field)后触发 _notify_listeners()。这个机制是 ReactiveRunner 事件驱动循环的关键:Runner 将自己注册为 listener,当节点写入状态时,Runner 被唤醒并驱动下一轮 tick。
sequenceDiagram
participant Node as Jianmu 节点
participant SM as StateManager
participant Runner as ReactiveRunner
participant Tree as 行为树
Node->>SM: write_state({"done": False})
SM->>SM: 验证 & 写入
SM->>Runner: _on_wake_signal()
Runner->>Runner: _signal_tick() → tick_signal.set()
Runner->>Tree: tree.tick()
Tree->>Node: update() / update_async()
Runner 的 _on_wake_signal 只在 auto_driving=True(即处于 run() 模式)时才会实际触发 tick 信号。在 step() 模式下,auto_driving=False,写入不会自动触发下一轮——由调用方显式控制节奏。
端到端流程:从 Schema 定义到 run() 执行¶
将以上所有机制串联起来,一个完整的 Agent 运行周期如下:
sequenceDiagram
participant User as 用户代码
participant Agent as Agent 外观
participant Runner as ReactiveRunner
participant SM as StateManager
participant Node as LLM/Skill 节点
User->>Agent: Agent(root, state_schema=MyState)
Agent->>SM: StateManager(MyState)
SM->>SM: _parse_schema() → 提取 reducer + ephemeral
Agent->>SM: initialize()
Agent->>Runner: ReactiveRunner(root, sm, ctx)
User->>Agent: agent.run({"messages": [human("hello")]})
Agent->>Runner: run(input_data)
Runner->>SM: reset_run_fields() → done=False, final_answer=""
Runner->>SM: update(input_data)
Runner->>Runner: 注入 deps + 绑定 callbacks
Runner->>Runner: auto_driving=True, _signal_tick()
loop 事件循环 tick
Runner->>Tree: tree.tick()
Tree->>Node: update_async()
Node->>SM: reset_call_fields()
Node->>SM: update/merge_dict/transform_field
SM->>Runner: _on_wake_signal() → 下一轮 tick
end
Runner->>SM: dump_checkpoint() → 保存
Runner-->>Agent: Status.SUCCESS
Agent-->>User: 执行完成
最佳实践与反模式¶
✅ 推荐做法:
- 为
messages类字段使用Annotated[List[Message], operator.add],保证追加语义而非覆盖 - 将 UI 关心的临时输出(如当前轮次的工具调用)标记为
Ephemeral("step") - 将 run 级别的输出(如
done、final_answer)标记为Ephemeral("run") - 在多节点并发写入同一字典的场景使用
merge_dict而非update - 在对列表/计数器执行增量操作时使用
transform_field以确保原子性
❌ 反模式:
- 在 namespace 写入中使用与 schema 字段同名的键——这会被硬性拒绝
- 在
extra='ignore'或extra='forbid'的 Schema 上使用 namespace 写入——会静默丢失或抛出异常 - 在 Reducer 中执行副作用——Reducer 应该是纯函数
- 依赖 Ephemeral 字段在
run()之间保持值——它们会被自动复位
下一步:了解了 StateManager 如何管理类型化状态后,建议继续阅读 ReactiveRunner:事件驱动的异步 tick 调度与挂起恢复机制,理解 StateManager 的 listener 通知如何驱动行为树的异步执行循环。