系列文章:
手把手构建企业级 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 交替生成“思考”和“行动”,并根据行动结果调整后续思考。整个循环可以抽象为以下状态机:

每一次循环都是一次有状态的交互:
- Thought(思考):LLM 读取当前对话历史、技能说明和可用工具列表,生成一段推理文本,并决定下一步动作。
- Action(行动):如果 LLM 输出了工具调用指令(例如
{"action": "tool_call", "name": "query_order", "params": {"order_id": "O123"}}),运行时解析并执行对应的工具。 - Observation(观察):工具执行的结果(成功返回值或异常信息)被封装后加入到对话历史中。
- 重复:运行时再次将更新后的历史交给 LLM,让它决定继续行动还是生成最终回复。
💡 关键设计决策: 我们将使用大模型原生的 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 pyyaml5.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 协议,并为企业级安全、多模型切换预留了扩展点。
九、总结与下一步
本文我们:
- 深入理解了 ReAct 循环的原理与状态机设计。
- 实现了基于 OpenAI function calling 的完整 AgentRuntime,支持工具调用、记忆读取和流式事件。
- 给出了 Claude tool use 的适配方案,确保框架的模型无关性。
- 编写了测试脚本,验证了 Agent 闭环。
下一篇文章预告:《手把手构建企业级 Agent 框架:Skill 系统--知识注入与能力扩展》。敬请期待!







