Token导航 LogoToken导航TokenDH.com

手把手构建企业级 Agent 框架(三):Pi Agent 运行时与 ReAct 循环

更新时间 2026-05-20来源 Aike正文 1.2万字阅读约 39分钟2 张图片
上次,我们实现了一个强大的 Gateway 网关,让消息能从飞书、WebChat 等不同渠道汇聚并安全路由。今天,我们要深入框架的灵魂——Agent 运行时(Pi Agent Runtime)。这正是让 Agent 从“你说我答”的对话机器进化为“你吩咐我执行”的自主代理的核心引擎。我们将亲手实现一个遵循 ReAct 模式的事件循环,并通过标准的 function calling / tool use 机制接入大模型,使 Agent 能够思考、决策、调用工具并从观察中学习。

系列文章:

手把手构建企业级 Agent 框架:从 OpenClaw 架构到自主实现

手把手构建企业级 Agent 框架(二):Gateway 网关与多渠道接入

一、Pi Agent 的定位:嵌入式引擎而非外挂框架

OpenClaw 中的 Pi Agent 并不是一个庞大的编排器,而是一个轻量级嵌入式运行时。它直接运行在消息回路中,负责协调 LLM 推理与工具执行。这种设计哲学意味着:

PS:Pi Agent 是一个遵循“极简主义”设计哲学的微型智能编程体(Coding Agent),它运行在本地设备(如 Mac/Linux 或树莓派)上,是 OpenClaw 架构中真正执行任务的“本地肢体”。在 OpenClaw 的整体架构里,Pi Agent 是关键的 Pi-embedded 执行端,与云端大脑(Orchestrator)、协议网关(Gateway)各司其职。这种设计使得 Pi Agent 能安全地运行在本地沙箱环境中,执行命令或调用技能,最终将结果反馈给用户。

  • 极低的抽象泄漏:
    开发者可以直接控制提示词、工具选择逻辑和错误恢复。
  • 流式原生:
    每一步思考、每一次工具调用都可以实时推送给用户,提供透明的执行体验。
  • 自包含:
    Agent 运行时可以独立测试,无需依赖 Gateway 或外部消息系统。

我们的企业级 Agent 运行时将完全继承这些优点,同时注入多模型容错、审计日志、安全上下文传递等企业刚需。

二、ReAct 循环详解:Thought → Action → Observation 的无限流

ReAct(Reasoning and Acting)是目前最主流的 Agent 推理范式。它让 LLM 交替生成“思考”和“行动”,并根据行动结果调整后续思考。整个循环可以抽象为以下状态机:

图片

每一次循环都是一次有状态的交互

  1. Thought(思考):
    LLM 读取当前对话历史、技能说明和可用工具列表,生成一段推理文本,并决定下一步动作。
  2. Action(行动):
    如果 LLM 输出了工具调用指令(例如 {"action": "tool_call", "name": "query_order", "params": {"order_id": "O123"}}),运行时解析并执行对应的工具。
  3. Observation(观察):
    工具执行的结果(成功返回值或异常信息)被封装后加入到对话历史中。
  4. 重复:
    运行时再次将更新后的历史交给 LLM,让它决定继续行动还是生成最终回复。
  5. 💡 关键设计决策: 我们将使用大模型原生的 function calling / tool use 能力,而非通过提示词工程让模型输出自定义 JSON。这能大幅提高工具选择的准确率,减少解析错误,并支持并行工具调用。下面会同时给出 OpenAI 和 Claude 两种实现。

三、Agent 运行时架构:事件驱动的异步循环

我们的实现采用 asyncio 异步生成器 作为核心事件流。这样既能支持流式输出(用户可实时看到思考过程),又能与 Gateway 的 WebSocket 完美结合。

核心组件架构图:

图片

AgentRuntime 将所有逻辑封装在 run() 异步生成器中,每次 yield 一个事件字典。Gateway 可以这样使用:

asyncfor event in agent.run(session_id, user_message):
# event: {"type": "thought", "data": "我需要查询..."}
# 或者 {"type": "tool_call", "data": {...}}
await websocket.send_json(event)

四、数据流与接口设计

我们定义一套清晰的数据结构来保证各模块解耦:

4.1 用户请求与上下文

from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional

@dataclass
classAgentRequest:
    session_id:str
    user_id:str
    tenant_id:str
    content:str
    metadata: Dict[str, Any]= field(default_factory=dict)

@dataclass
classAgentContext:
"""贯穿整个 Agent 循环的安全与配置上下文"""
    request: AgentRequest
    security_roles: List[str]= field(default_factory=list)
    trace_id: Optional[str]=None

4.2 事件流定义

classAgentEventType:
    THOUGHT ="thought"# LLM 推理文本
    TOOL_CALL ="tool_call"# 即将调用工具
    TOOL_RESULT ="tool_result"# 工具执行结果
    TEXT_DELTA ="text_delta"# 最终回复的流式片段
    FINAL ="final"# 回复完成,可附带完整文本
    ERROR ="error"# 异常事件

4.3 工具接口

工具定义需要遵循 OpenAI function calling 的 JSON Schema 规范,这也是大多数模型通用的工具描述格式。

from dataclasses import dataclass
from typing import Any, Callable, Dict, List

@dataclass
classToolDefinition:
    name:str
    description:str
    parameters: Dict  # JSON Schema,如 {"type":"object","properties":{...},"required":[...]}
    fn: Callable      # 异步可调用对象

classToolRegistry:
def__init__(self):
        self._tools: Dict[str, ToolDefinition]={}
defregister(self, td: ToolDefinition):...
defget_tool_schemas(self)-> List[Dict]:...# 供 LLM 使用的工具列表
asyncdefexecute(self, name:str, params: Dict, ctx: AgentContext)-> Any:...

五、核心代码实战:基于 Function Calling 的 ReAct Agent 运行时

下面是我们今天的主角——AgentRuntime 的完整实现。核心变化:不再要求模型输出自定义 JSON,而是使用 OpenAI 标准的 function calling 能力。同时给出 Claude tool use 的适配代码,二者接口高度统一。

5.1 安装依赖

pip install openai anthropic fastapi uvicorn websockets pyyaml

5.2 项目结构

eclaw-agent/
├── main.py          # 启动 Agent 服务(可选)
├── agent.py         # AgentRuntime 核心(使用 function calling)
├── tools.py         # 工具注册与示例
├── memory.py        # 短期记忆(复用第一篇)
├── skill.py         # Skill 加载器(复用第一篇)
└── test_agent.py    # 独立测试脚本

5.3 Agent 运行时核心 (agent.py) —— 采用 OpenAI Function Calling

下面的实现严格遵循 OpenAI 的 function calling 协议:将工具定义传入 tools 参数,模型会返回 tool_calls 数组;我们执行工具后将结果以 role: "tool" 消息回传。

import json
import asyncio
from typing import AsyncIterator, Dict, List, Any
from openai import AsyncOpenAI

from tools import ToolRegistry, ToolDefinition
from memory import MemoryStore
from skill import SkillLoader
from agent_context import AgentContext

classAgentRuntime:
"""基于 OpenAI function calling 的 ReAct Agent 运行时"""
def__init__(self, max_steps:int=5, model_name:str="gpt-4o-mini"):
        self.max_steps = max_steps
        self.model_name = model_name
        self.tools = ToolRegistry()
        self.tools.register_default_tools()
        self.memory = MemoryStore()
        self.skills = SkillLoader()
# 初始化 OpenAI 客户端(兼容任何 OpenAI API 兼容服务)
        self.client = AsyncOpenAI()

asyncdefrun(self, session_id:str, user_input:str, ctx: AgentContext =None)-> AsyncIterator[dict]:
"""主循环:不断调用 LLM,直到生成最终回复或达到最大步数。"""
# 1. 加载短期记忆
        history =await self.memory.get_history(session_id)
# 2. 获取相关技能
        skill_text =await self.skills.load_relevant(user_input)
# 3. 构造消息列表
        messages =[
{"role":"system","content": self._build_system_prompt(skill_text)},
]
if history:
            messages.append({"role":"system","content":f"历史对话摘要:{history}"})
        messages.append({"role":"user","content": user_input})

        step =0
while step < self.max_steps:
            step +=1
# 调用 LLM(带 function calling)
            response =await self._call_llm(messages)
if response.get("error"):
yield{"type":"error","data": response["error"]}
return

            choice = response["choices"][0]
            msg = choice["message"]
            finish_reason = choice.get("finish_reason")

# 如果模型直接返回文本且没有工具调用,即为最终回复
if finish_reason =="stop"andnot msg.get("tool_calls"):
                answer = msg.get("content","")
yield{"type":"final","data": answer}
await self.memory.add(session_id,"user", user_input)
await self.memory.add(session_id,"assistant", answer)
return

# 处理 tool_calls
if msg.get("tool_calls"):
# 将 assistant 消息(含 tool_calls)加入历史
                messages.append(msg)
for tool_call in msg["tool_calls"]:
                    tool_name = tool_call["function"]["name"]
                    tool_params = json.loads(tool_call["function"]["arguments"])
yield{"type":"tool_call","data":{"tool": tool_name,"params": tool_params}}

# 执行工具
try:
                        tool_result =await self.tools.execute(tool_name, tool_params, ctx)
yield{"type":"tool_result","data": tool_result}
                        result_str = json.dumps(tool_result, ensure_ascii=False)
except Exception as e:
                        error_str =f"工具执行失败: {str(e)}"
yield{"type":"tool_result","data":{"error": error_str}}
                        result_str = error_str

# 将工具结果以 tool 角色加入对话
                    messages.append({
"role":"tool",
"tool_call_id": tool_call["id"],
"content": result_str
})
else:
# 如果没有 tool_calls 但也没 stop,可能是其他情况,添加模型消息继续循环
                messages.append(msg)

yield{"type":"final","data":"抱歉,任务似乎太复杂了,请简化请求或联系管理员。"}

def_build_system_prompt(self, skill_text:str)->str:
"""构造系统提示词,不再要求模型输出自定义 JSON"""
return(
"你是一个企业级智能助手,能够使用工具完成任务。\n"
f"相关业务技能指导:\n{skill_text}\n"
"请根据用户需求,自主决定是否需要调用工具。"
)

asyncdef_call_llm(self, messages: List[Dict])-> Dict:
"""调用 OpenAI 接口,传入工具定义"""
        tool_schemas = self.tools.get_tool_schemas()# 符合 OpenAI 的 tools 格式
try:
            response =await self.client.chat.completions.create(
                model=self.model_name,
                messages=messages,
                tools=tool_schemas if tool_schemas elseNone,
                tool_choice="auto"if tool_schemas elseNone,
                temperature=0.2,
)
return response.model_dump()
except Exception as e:
return{"error":f"LLM 调用失败: {str(e)}"}

5.4 工具模块 (tools.py) —— 符合 OpenAI 工具格式

import asyncio
from typing import Dict, List, Any
from dataclasses import dataclass

@dataclass
classToolDefinition:
    name:str
    description:str
    parameters: Dict  # JSON Schema
    fn: Any

classToolRegistry:
def__init__(self):
        self._tools: Dict[str, ToolDefinition]={}

defregister(self, td: ToolDefinition):
        self._tools[td.name]= td

defget_tool_schemas(self)-> List[Dict]:
"""返回 OpenAI 兼容的工具列表"""
return[
{
"type":"function",
"function":{
"name": td.name,
"description": td.description,
"parameters": td.parameters,
}
}
for td in self._tools.values()
]

asyncdefexecute(self, name:str, params: Dict, ctx)-> Any:
        td = self._tools.get(name)
ifnot td:
raise ValueError(f"工具 '{name}' 未注册")
if asyncio.iscoroutinefunction(td.fn):
returnawait td.fn(params, ctx)
else:
return td.fn(params, ctx)

defregister_default_tools(self):
asyncdefquery_order(params, ctx):
return{"order_id": params.get("order_id"),"status":"已发货","eta":"2026-05-20"}
        self.register(ToolDefinition(
            name="query_order",
            description="查询订单状态,需要提供 order_id",
            parameters={
"type":"object",
"properties":{
"order_id":{"type":"string","description":"订单号"}
},
"required":["order_id"]
},
            fn=query_order
))

asyncdefsend_notification(params, ctx):
return{"success":True,"message":f"已向 {params.get('user')} 发送通知"}
        self.register(ToolDefinition(
            name="send_notification",
            description="向指定用户发送通知",
            parameters={
"type":"object",
"properties":{
"user":{"type":"string","description":"接收用户ID"},
"content":{"type":"string","description":"通知内容"}
},
"required":["user","content"]
},
            fn=send_notification
))

5.5 适配 Claude 的 Tool Use 变体

如果使用 Anthropic Claude,只需修改 _call_llm 方法和工具格式。以下是完整对照代码,可直接替换上述 AgentRuntime 中的 LLM 调用部分:

# 使用 Claude Messages API 的版本
from anthropic import AsyncAnthropic

classClaudeAgentRuntime(AgentRuntime):
"""基于 Anthropic Claude tool use 的运行时"""
def__init__(self, max_steps=5, model_name="claude-sonnet-4-20250514"):
super().__init__(max_steps, model_name)
        self.client = AsyncAnthropic()

defget_tool_schemas(self)-> List[Dict]:
"""返回 Anthropic 兼容的工具格式(与 OpenAI 略有不同)"""
return[
{
"name": td.name,
"description": td.description,
"input_schema": td.parameters  # Anthropic 使用 input_schema
}
for td in self.tools._tools.values()
]

asyncdef_call_llm(self, messages: List[Dict])-> Dict:
# Anthropic 的系统提示词需要单独传递
        system_msg =""
        api_messages =[]
for m in messages:
if m["role"]=="system":
                system_msg += m["content"]+"\n"
else:
                api_messages.append(m)
        tool_schemas = self.get_tool_schemas()
try:
            response =await self.client.messages.create(
                model=self.model_name,
                system=system_msg.strip(),
                messages=api_messages,
                tools=tool_schemas if tool_schemas elseNone,
                max_tokens=1024,
)
# 转换为类似 OpenAI 的统一格式
return self._normalize_claude_response(response)
except Exception as e:
return{"error":f"Claude 调用失败: {str(e)}"}

def_normalize_claude_response(self, response)-> Dict:
"""将 Claude 响应转为统一格式以便循环处理"""
        text_content =""
        tool_calls =[]
for block in response.content:
if block.type=="text":
                text_content += block.text
elif block.type=="tool_use":
                tool_calls.append({
"id": block.id,
"function":{
"name": block.name,
"arguments": json.dumps(block.input)
}
})
return{
"choices":[{
"message":{
"role":"assistant",
"content": text_content,
"tool_calls": tool_calls
},
"finish_reason":"tool_calls"if tool_calls else"stop"
}]
}

通过统一的内部格式,Agent 的主循环无需任何改动即可在 OpenAI 和 Claude 之间切换。

5.6 短期记忆与 Skill 加载器(复用)

# memory.py
from collections import defaultdict

classMemoryStore:
def__init__(self, max_len=10):
        self.store = defaultdict(list)
        self.max_len = max_len
asyncdefadd(self, session_id, role, content):
        self.store[session_id].append({"role": role,"content": content})
iflen(self.store[session_id])> self.max_len *2:
            self.store[session_id]= self.store[session_id][-self.max_len*2:]
asyncdefget_history(self, session_id):
return"\n".join([f"{m['role']}: {m['content']}"for m in self.store.get(session_id,[])])
# skill.py
from pathlib import Path

classSkillLoader:
def__init__(self, skill_dir="skills"):
        self.skill_dir = Path(skill_dir)
        self.skills ={}
        self._load_all()
def_load_all(self):
ifnot self.skill_dir.exists():
            self.skill_dir.mkdir()
(self.skill_dir /"order_skill.md").write_text("# 订单技能\n用于处理订单查询。调用 query_order 工具时需提供 order_id。\n")
forfilein self.skill_dir.glob("*.md"):
            self.skills[file.stem]=file.read_text(encoding="utf-8")
asyncdefload_relevant(self, query:str)->str:
        relevant =[]
for name, content in self.skills.items():
ifany(word in query for word in name.split('_')):
                relevant.append(content)
return"\n".join(relevant)if relevant else"无特定技能"

六、运行与测试

我们编写一个测试脚本,模拟会话并观察 Agent 基于 function calling 的推理过程。

6.1 测试脚本 (test_agent.py)

import asyncio
from agent import AgentRuntime
from agent_context import AgentContext

asyncdefmain():
    agent = AgentRuntime()
    ctx = AgentContext(request=None, security_roles=["user"])
    session_id ="test-session-1"

print("=== 测试1:查询订单 ===")
asyncfor event in agent.run(session_id,"帮我查一下订单 O12345 的状态", ctx):
print(f"[{event['type']}] {event['data']}")

print("\n=== 测试2:简单问候 ===")
asyncfor event in agent.run(session_id,"你好啊", ctx):
print(f"[{event['type']}] {event['data']}")

if __name__ =="__main__":
    asyncio.run(main())

6.2 预期输出

=== 测试1:查询订单 ===
[tool_call] {'tool': 'query_order', 'params': {'order_id': 'O12345'}}
[tool_result] {'order_id': 'O12345', 'status': '已发货', 'eta': '2026-05-20'}
[final] 您的订单 O12345 目前状态为“已发货”,预计送达时间 2026-05-20。

注意:由于使用了真正的 function calling,模型返回的最终回复文本会更加自然,不再需要我们硬编码格式。

📌 实际接入真实 API 的注意事项:

  • 确保 AsyncOpenAI 的 API key 已配置,或指向兼容服务(如 Azure、vLLM)。
  • 工具描述的 JSON Schema 尽量详细,这直接影响模型的选择准确性。
  • 企业环境中应加入请求缓存,避免重复调用 LLM。
  • 可在 tool_choice 参数中设置 "required" 强制模型调用工具。

七、流式处理与高级容错

上面的代码已展示了基本事件流。在企业级场景,我们还需要:

  • 真正的流式文本:
    OpenAI 支持 stream=True,可实时推送 text_delta 事件。
  • 并行工具调用:
    OpenAI 支持一次返回多个 tool_calls,用 asyncio.gather 并行执行。
  • 智能重试与模型回退:
    LLM 调用失败自动重试,主模型不可用时切换备用模型。
  • 工具超时与熔断:
    通过 asyncio.wait_for 控制执行时间,连续失败时自动熔断。

八、与 OpenClaw 的对标思考

🔍 “我们做了什么” vs “OpenClaw 为什么这样做”?

Function Calling: OpenClaw 的 Pi Agent 使用 TypeScript 实现,通过提示词引导模型输出结构化 JSON。我们的方案直接使用大模型原生的 function calling / tool use 能力,准确率更高、解析更可靠,且天然支持并行工具调用。这是企业级 Agent 的推荐实践。

模型兼容性: 通过将 OpenAI 和 Claude 的调用封装为统一接口,我们的运行时可以灵活切换底层模型,而不影响上层 ReAct 循环。OpenClaw 目前主要适配 OpenAI 生态。

安全性: 我们的 AgentContext 贯穿整个执行链,工具执行前可校验权限;OpenClaw 的个人助手模式默认信任本地用户。在企业改造时,安全上下文显式传递至关重要。

总而言之,我们的 Agent 运行时在保持 OpenClaw 简洁高效基因的同时,采用了更符合行业标准的 function calling 协议,并为企业级安全、多模型切换预留了扩展点。

九、总结与下一步

本文我们:

  1. 深入理解了 ReAct 循环的原理与状态机设计。
  2. 实现了基于 OpenAI function calling 的完整 AgentRuntime,支持工具调用、记忆读取和流式事件。
  3. 给出了 Claude tool use 的适配方案,确保框架的模型无关性。
  4. 编写了测试脚本,验证了 Agent 闭环。

下一篇文章预告:《手把手构建企业级 Agent 框架:Skill 系统--知识注入与能力扩展》。敬请期待!

文章标签智能体
资讯来源:由AI资讯编辑整理自互联网公开内容,版权归原作者所有,未经许可,不得转载。

继续浏览更多资讯

返回资讯目录

相关资讯

更多