数据截至 (上游 commit 1f738cdeb7f5)
Agent 抽象层:run 循环、ChatClient 与会话
30 秒导读: 在 Agent Framework 里,一个「agent」不是什么魔法,它是把三样东西包成一个
run()入口的薄壳:一个 ChatClient(负责真正调模型)、一套 统一的消息/响应模型(不管底层是 OpenAI 还是 Anthropic,进出都是Message/ChatResponse)、一个 会话(决定对话历史存在本地内存还是服务端)。本章只讲「单个 agent 怎么跑一圈」;工具调用循环和 middleware 留给 第 02 章。
本章覆盖 python/packages/core/agent_framework/ 下四个核心文件:_agents.py(agent 本体)、_clients.py(chat client 抽象)、_types.py(消息与响应模型)、_sessions.py(会话与历史)。
1. 这是什么(零基础也能懂)
一句话定义: 一个 Agent = 一个聊天客户端(ChatClient)+ 一层统一的消息/响应封装 + 一个可选的会话(记住上下文)。
解决什么问题: 假设你想让程序「和一个大模型对话、还能记住上下文、还能换供应商不改业务代码」。裸调 OpenAI SDK 做不到这些——每家 API 的请求/响应格式不同,历史管理各有各的规矩。Agent Framework 把这些差异抹平:你面对的永远是同一个 agent.run("...")。
用起来什么样: 最小的一次调用长这样。
# 示意,非源码
from agent_framework.openai import OpenAIChatClient
# 1. 造一个 chat client(负责真正连模型)
client = OpenAIChatClient(model="gpt-4o")
# 2. 用它造一个 agent,给它名字和系统指令
agent = client.as_agent(name="assistant", instructions="你是一个乐于助人的助手。")
# 3. 跑一圈,拿最终回复
response = await agent.run("巴黎天气怎么样?")
print(response.text) # 完整回复文本
# 4. 想要边生成边显示?同一个 run,加 stream=True
async for update in agent.run("再讲个笑话", stream=True):
print(update.text, end="")
注意第 2 步的 client.as_agent(...)——这是 BaseChatClient.as_agent()(python/packages/core/agent_framework/_clients.py:571)提供的便捷方法,等价于 Agent(client=client, ...)。
一句话直觉: 把 ChatClient 想成「一根接到某家大模型的电话线」,把 Agent 想成「拿着这根电话线、还带了记事本(会话)和话术卡(instructions)的接线员」。你只跟接线员说话,不用管电话线接到哪家。
本节不出现底层细节。记住一件事:agent 自己几乎不「思考」,它把活儿委托给 chat client,自己只管统一入口、统一格式、管理会话。
2. 顶层全景(它大概怎么转)
2.1 一次 run() 的数据流
先看「一次非流式 run() 从入口到出口」的主线。怎么读这张图:从上到下是时间顺序,左边是 agent 层做的事,右边标注落在哪个文件。
agent.run("巴黎天气?") 文件 / 符号
│
▼
① 组装运行上下文 _agents.py:1160
· normalize_messages 把输入 _prepare_run_context
统一成 list[Message]
· 合并 default_options + 运行期 options (_merge_options)
· 决定用本地历史还是服务端历史
· 跑 before_run 加载会话历史
│
▼
② 调下游 chat client _agents.py:995
client.get_response(messages, ...) _call_chat_client
│
▼
③ chat client 真正连模型 _clients.py:479 get_response
· 具体供应商实现 _inner_get_response _clients.py:412 (抽象)
· 返回统一的 ChatResponse
│
▼
④ 收尾 _agents.py:1023
· 给每条消息盖上 author_name _parse_non_streaming_response
· conversation_id 回写会话 _finalize_response (:1377)
· 跑 after_run 存历史
│
▼
AgentResponse(messages=..., text=...)
流式路径(stream=True)走的是同一个 run(),只是在 ② 之后返回一个 ResponseStream,把上面 ④ 的收尾挂成流结束后的钩子(见 §3.2)。
2.2 主要部件一句话 职责
| 部件 | 干什么 | 在哪(文件:符号) |
|---|---|---|
SupportsAgentRun | agent 的「接口契约」(协议),定义 run / create_session 等 | _agents.py:234 |
BaseAgent | 最小 agent 基类:管 id/name、context providers、会话 | _agents.py:374 |
RawAgent | 核心实现:把 chat client 包起来,真正跑 run 循环 | _agents.py:730 |
Agent | 生产用的类:RawAgent + middleware 层 + 遥测层 | _agents.py:1786 |
SupportsChatGetResponse | chat client 的接口契约 | _clients.py:85 |
BaseChatClient | chat client 抽象基类,供应商子类继承它 | _clients.py:217 |
Message / ChatResponse / AgentResponse | 统一的消息与响应模型 | _types.py:1678 / 2128 / 2536 |
ResponseStream | 流式包装:能迭代、也能在最后聚合成完整响应 | _types.py:2939 |
AgentSession | 轻量会话状态容器(session id + state) | _sessions.py:1717 |
HistoryProvider | 历史存储策略基类(内存、文件…) | _sessions.py:943 |
2.3 四层类的继承关系
Agent 不是一个巨类,而是叠罗汉叠出来的。怎么读:越往下功能越全,Agent 在最底下、能力最全。
BaseAgent 最小骨架:id/name/会话/context providers,不能直接 run
│ (被继承)
▼
RawAgent 核心 run 循环:包 chat client、统一消息/响应、流式非流式
│ (被继承,再叠两层 mixin)
▼
Agent = AgentMiddlewareLayer + AgentTelemetryLayer + RawAgent
└── 加 middleware 拦截 └── 加 OpenTelemetry 观测
依据:class Agent(AgentMiddlewareLayer, AgentTelemetryLayer, RawAgent[OptionsCoT], ...)(_agents.py:1786)。Agent.run() 只是把参数原样转给 super().run()(_agents.py:1849),真正的循环在 RawAgent。middleware/telemetry 两层留给第 02 章。
3. 核心原理(逐个机制)
3.1 统一的消息 / 响应模型
要解决的小问题: 不同供应商的消息格式五花八门。框架内部必须有一套自己的、供应商无关的类型,进出都用它。
四个主角:
| 类型 | 是什么 | 关键属性 |
|---|---|---|
Message | 一条消息 | role、contents(内容项列表)、author_name、.text(拼接文本) |
ChatResponse | chat client 层的完整响应 | messages、response_id、conversation_id、usage_details |
AgentResponse | agent 层的完整响应 | messages、response_id、agent_id、.text |
*ResponseUpdate | 流式的单个增量块 | contents、role、message_id… |
为什么 chat client 层和 agent 层各有一套(ChatResponse vs AgentResponse)?因为职责不同:ChatResponse 贴近「一次模型调用」的原始结果,AgentResponse 贴近「一次 agent run」的对外结果(带 agent_id,可能聚合了多次模型调用)。
入口先统一:normalize_messages。 你可以给 run() 传字符串、传 Content、传 Message、传一个混合列表——框架第一步就把它们全捏成 list[Message]。
# 示意,非源码。真实逻辑见 _types.py:1774 normalize_messages
def normalize_messages(messages):
if messages is None:
return []
if isinstance(messages, str):
return [Message("user", [messages])] # 裸字符串 => 一条 user 消息
if isinstance(messages, Message):
return [messages]
# 列表:逐个规整,字符串/Content 都包成 user 消息
return [m if isinstance(m, Message) else Message("user", [m]) for m in messages]
真实实现见 normalize_messages(_types.py:1774),输入类型别名 AgentRunInputs(_types.py:1771)。Message.text 属性把所有文本内容项拼起来(_types.py:1761)。
流式增量怎么变成 agent 的增量:map_chat_to_agent_update。 chat client 吐出的是 ChatResponseUpdate,agent 要对外给 AgentResponseUpdate。转换是一对一的字段搬运,顺手把 author_name 补上 agent 名字:
# 真实源码,_types.py:2917 map_chat_to_agent_update
def map_chat_to_agent_update(update: ChatResponseUpdate, agent_name: str | None):
return AgentResponseUpdate(
contents=update.contents,
author_name=update.author_name or agent_name, # 没作者名就填 agent 名
response_id=update.response_id,
message_id=update.message_id,
raw_representation=update, # 原始块留底
...
)
增量聚合成完整响应:from_updates。 一堆 update 怎么拼回一个 ChatResponse / AgentResponse?靠 _process_update(_types.py:1869)判断每个 update 是延续当前消息还是开一条新消息(看 message_id),再由 ChatResponse.from_updates(_types.py:2284)/ AgentResponse.from_updates 汇总。这就是流式结束时「碎片拼整」的地方。
3.2 同一个 run() 入口:流式与非流式
要解决的小问题: 大多数框架把「一次拿完」和「边生成边收」拆成两个方法(如 get_response / get_streaming_response)。Agent Framework 只用一个 run(),靠 stream: bool 参数分流,并用**函数重载(overload)**让类型检查器知道返回什么。
同一个签名,三种重载:
| 你怎么调 | 静态返回类型 | 依据 |
|---|---|---|
run(msg) 或 run(msg, stream=False) | Awaitable[AgentResponse] | _agents.py:1006、:869 |
run(msg, stream=True) | ResponseStream[AgentResponseUpdate, AgentResponse] | _agents.py:1036 |
重载只是给 IDE/类型检查器看的「门面」;真正的实现是最后那个 def run(...)(_agents.py:1051),内部按 stream 分岔:
# 示意,非源码。骨架取自 _agents.py:960-977
def run(self, messages=None, *, stream=False, ...):
if not stream:
async def _run_non_streaming():
ctx = await self._prepare_run_context(...) # 组装上下文
response = await self._call_chat_client(ctx, stream=False)
return await self._parse_non_streaming_response(ctx, response)
return _run_non_streaming() # 返回一个协程
async def _run_streaming():
ctx = await self._prepare_run_context(...)
stream_response = self._call_chat_client(ctx, stream=True)
return self._parse_streaming_response(ctx, stream_response)
return ResponseStream.from_awaitable(_run_streaming()) # 返回一个流
重点看:两条路共享 _prepare_run_context 和 _call_chat_client。 区别只在最后的收尾——非流式直接 await 拿 ChatResponse 再转成 AgentResponse;流式返回 ResponseStream,把收尾逻辑挂成钩子。
ResponseStream 是什么(_types.py:2939): 一个既能 async for 迭代、又能在最后 await stream.get_final_response() 拿完整响应的包装器。它的巧处在 map(_types.py:2990):把内层 ChatResponseUpdate 流实时转成 AgentResponseUpdate 流,同时保证内层的收尾钩子(存历史、遥测)在最后仍然执行。
# 真实源码骨架,_agents.py:1112 _parse_streaming_response
stream = stream_response.map(
transform=partial(map_chat_to_agent_update, agent_name=self.name), # 每块转换
finalizer=_finalizer, # 结束时聚合
)
return stream.with_transform_hook(...).with_result_hook(_post_hook) # 挂收尾钩子
_post_hook(_agents.py:1253)在流跑完后回写 conversation_id 到会话、给消息补作者名、跑 after_run 存历史——正好对应 §2.1 图里非流式那一步的 ④,只是延后到流结束。
3.3 Agent 如何「包」一个 chat client
要解决的小问题: agent 构造时收下一个 client,还有一堆默认选项(温度、工具、指令…);运行时又可能传另一套选项。谁覆盖谁?怎么合?
契约:SupportsChatGetResponse(_clients.py:85)。 agent 不依赖具体 client 类,只依赖这个协议——只要有个符合签名的 get_response,就能当 agent 的引擎。协议用结构化子类型(duck typing),client 不需要显式继承。
构造时:default_options 一次性拍平。 RawAgent.__init__(_agents.py:811)把 instructions、tools、temperature 等全塞进一个 self.default_options 字典(_agents.py:928),并剔除所有 None 值(_agents.py:928)。MCP 工具单独存进 self.mcp_tools,运行时才展开成函数(_agents.py:899)。
运行时:_merge_options 决定覆盖规则。 运行期选项和 agent 默认选项在 _prepare_run_context(_agents.py:1353)里合并,核心是 _merge_options(_agents.py:174)。规则不是简单覆盖:
| 选项 | 合并方式 | 依据 |
|---|---|---|
tools | 去重后追加(名字冲突报错) | _agents.py:157 |
logit_bias / metadata | 两个字典合并 | _agents.py:165-170 |
instructions | 字符串拼接(换行连起) | _agents.py:171 |
| 其它 | 运行期值覆盖默认值 | _agents.py:174 |
值为 None | 跳过(视为「未设置」,不覆盖) | _agents.py:154 |
最后再扫一遍剔掉所有 None(_agents.py:176)——这样「没设的选项」不会被硬塞给服务,交给服务用它自己的默认。
一个细节:agent 名字要能当函数名。 当 agent 被当成工具用(as_tool,_agents.py:608-613)或暴露成 MCP server 时,名字要合法。_sanitize_agent_name(_agents.py:179)把空格/特殊字符换成下划线、压缩连续下划线、开头是数字就加前缀,空了就兜底成 "agent"。
_RunContext(_agents.py:216): 上面所有准备工作的产物打成一个 TypedDict——session、session_messages、agent_name、chat_options、suppress_response_id 等,一路传给调用和收尾环节。它就是「这一次 run 的全部上下文」。
3.4 provider 无关抽象与能力协议
要解决的小问题: 框架怎么在「所有 client 长一样」和「有些 client 会特殊技能(代码解释器、联网搜索…)」之间取平衡?
一半靠继承,一半靠协议。
BaseChatClient(_clients.py:217)是所有 client 的抽象基类。它把公共流程写死在 get_response(_clients.py:482)里:解析 compaction 覆盖、合并 client_kwargs,然后调一个抽象方法 _inner_get_response(_clients.py:415)。供应商子类只需实现 _inner_get_response,同时处理流式和非流式(看 stream 参数),别的都白拿。
# 示意,非源码。自定义 client 的最小形态,对应 _clients.py:240 的文档示例
class MyChatClient(BaseChatClient):
async def _inner_get_response(self, *, messages, stream, options, **kwargs):
if stream:
async def _gen():
yield ChatResponseUpdate(role="assistant", contents=[...])
return _gen()
return ChatResponse(messages=[Message("assistant", ["嗨!"])])
特殊能力用「能力协议」表达,而不是塞进基类。 _clients.py 定义了一组 Supports*Tool 协议,运行时用 isinstance 探测:
| 协议 | 代表「这个 client 会…」 | 依据 |
|---|---|---|
SupportsCodeInterpreterTool | 提供代码解释器工具 | _clients.py:668 |
SupportsWebSearchTool | 提供联网搜索 工具 | _clients.py:698 |
SupportsImageGenerationTool | 提供图像生成工具 | _clients.py:728 |
SupportsMCPTool | 提供 MCP 工具 | _clients.py:758 |
SupportsFileSearchTool | 提供文件检索工具 | _clients.py:789 |
SupportsShellTool | 提供 shell 执行工具 | _clients.py:819 |
用起来就是探测再取工具:
# 示意,非源码。模式取自 _clients.py:676 的文档
if isinstance(client, SupportsCodeInterpreterTool):
tool = client.get_code_interpreter_tool() # 只有会这招的 client 才有这方法
agent = client.as_agent(tools=[tool])
好处:基类保持精简,不会为了某家供应商的独门功能而膨胀;能不能用某功能,一个 isinstance 就知道。嵌入(embedding)也走同样的套路——SupportsGetEmbeddings 协议(_clients.py:871)+ BaseEmbeddingClient 基类(_clients.py:926)。
3.5 会话如何承载对话历史
要解决的小问题: 多轮对话要记住上文。但「记在哪」有两种截然不同的世界:本地(框架自己存消息列表)和服务端(像 OpenAI Responses API,服务器记着,你只拿一个 id)。会话层要同时罩住这两种。
AgentSession(_sessions.py:1717)是个轻量容器,只装三样:
| 字段 | 含义 |
|---|---|
session_id | 本地会话 id(不传就自动生成 UUID) |
service_session_id | 服务端会话 id(用服务端存储时才有) |
state | 一个可变字典,历史/上下文都塞这里 |
关键设计:provider 归 agent 所有,不归 session。 session 只存 id 和 state,真正「怎么存历史」的逻辑在 HistoryProvider。这样一个 session 可以在不同 provider 间流转(_sessions.py:1719 注释)。
两条历史世界线:
有 service_session_id / conversation_id?
或 client.STORES_BY_DEFAULT?
│
┌──────────────┴───────────────┐
否 是
│ │
┌─────────▼─────────┐ ┌──────────▼──────────┐
│ 本地历史 │ │ 服务端历史 │
│ 框架自己存消息 │ │ 服务器记着,框架 │
│ InMemoryHistory │ │ 只回写/带上那个 id │
│ Provider 自动挂上 │ │ │
└───────────────────┘ └─────────────────────┘
本地世界:自动挂 InMemoryHistoryProvider。 当你传了 session、又没配任何 context provider、又没有服务端存储迹象时,agent 自动给你塞一个内存历史 provider(_agents.py:1396-1405)。它把消息存在 session.state["messages"] 里(InMemoryHistoryProvider,_sessions.py:2087)。判断「服务端是否默认存储」看的是 client 的类属性 STORES_BY_DEFAULT(_clients.py:279)。
provider 的两个回调:before / after。 ContextProvider(_sessions.py:793)定义了 before_run(调模型前,往上下文里加历史/指令/工具)和 after_run(调模型后,把这轮消息存下来)。HistoryProvider(_sessions.py:943)是它的历史专用子类,只要实现 get_messages / save_messages 两个抽象方法,加载/存储的时机由基类的默认 before/after 处理。内置实现有内存版和文件版(FileHistoryProvider,一会话一个 JSONL 文件,_sessions.py:2171)。
本地 vs 服务端 id 的哨兵值。 服务端 id 会在 run 之后回写进 session.service_session_id(见 _update_session_from_chat_response,_agents.py:1185-1199;旧版在 _finalize_response 里)。但本地历史模式下需要一个「假」的 conversation id 让工具循环认得——这就是哨兵常量 LOCAL_HISTORY_CONVERSATION_ID(_sessions.py:1310)和判定函数 is_local_history_conversation_id(_sessions.py:1313)。框架用它区分「这是真服务端 id,要回写会话」还是「这是本地占位符,别当真」(_agents.py:1264-1270)。
每次模型调用都落一次历史:PerServiceCallHistoryPersistingMiddleware(_sessions.py:1535)。 默认情况下历史是「一次 run 存一次」。但如果开了 require_per_service_call_history_persistence,这个 middleware 会在每次模型调用前后加载/持久化历史。它有两种行为,由 service_stores_history 切换:
- 服务端不存:middleware 自己加载历史、把本地哨兵 id 塞进去驱动工具循环,调真 client 前再把哨兵剥掉(
_strip_local_conversation_id,_sessions.py:1613)。 - 服务端存:middleware 只负责「写一份到 provider」,真 conversation id 原样透传(
_sessions.py:1685-1693)。
这是「审计/评估要每步留痕」这类场景用的,细节和它如何嵌进 middleware 链在第 02 章。
4. 巧妙之处(可借鉴的技术)
-
一个
run()+stream: bool+ 重载,取代两个方法。 用户只学一个入口,类型检查器靠@overload仍能精确推断流式/非流式的返回类型(_agents.py:1006-1049)。少一半 API 表面积。 -
None= 「未设置」而非「设成空」。_merge_options(_agents.py:174)全程把None当「跳过」,最后统一剔除。好处:agent 默认值、运行期值、服务端默认值三层能干净叠加,没人会用一个None意外抹掉别人设好的值。 -
能力用协议探测,不用基类膨胀。
Supports*Tool一族(_clients.py:668起)让「会不会某招」变成一次isinstance,基类始终精简。加新能力不用改BaseChatClient。 -
本地/服务端历史用哨兵 id 统一。 一个
LOCAL_HISTORY_CONVERSATION_ID(_sessions.py:1310)让工具循环的代码不必到处if 本地 else 服务端——本地也有个「id」,只是框架内部认得、绝不外传给真 client。 -
流式的收尾钩子不会因
map丢失。ResponseStream.map(_types.py:2990)刻意保证内层 finalizer 和 result_hook 都执行,所以「存历史 / 回写会话 / 遥测」在流式下照样跑,只是延到流结束(_agents.py:1253_post_hook)。
5. 边界与局限
-
BaseAgent/BaseChatClient不能直接实例化。 前者没实现run(_agents.py:383注释),后者是 ABC 且_inner_get_response抽象(_clients.py:415)。要么用Agent/ 具体 client,要么自己实现协议。 -
不支持函数调用的 client 只会警告,不报错。 构造 agent 时如果 client 不是
FunctionInvocationLayer,只打一条 warning(_agents.py:872),工具能力会受限——不会拦你。 -
as_mcp_server只转发文本内容。 agent 暴露成 MCP server 时,工具结果里的图像/音频等富内容目前被丢弃并告警(_agents.py:1761),只有 text 会过 MCP。 -
服务端存储模式下,client 不回
conversation_id会静默丢跨轮历史。 框架每次都告警提醒(_sessions.py:1647),但不会替你补救。 -
require_per_service_call_history_persistence不能配已有的服务端会话。 本地哨兵驱动的工具循环无法和一个真实服务端会话对账,遇到会直接抛AgentInvalidRequestException(_agents.py:971)。
6. 代码地图(导航索引)
用符号名 grep 比行号抗漂移。下表是本章涉及的关 键落点。
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| agent 接口契约 | python/packages/core/agent_framework/_agents.py | SupportsAgentRun |
| 最小 agent 基类 | python/packages/core/agent_framework/_agents.py | BaseAgent |
| 核心 run 循环实现 | python/packages/core/agent_framework/_agents.py | RawAgent |
| 生产用 agent(叠 middleware/遥测) | python/packages/core/agent_framework/_agents.py | Agent |
| 组装单次 run 上下文 | python/packages/core/agent_framework/_agents.py | RawAgent._prepare_run_context |
| 选项合并规则 | python/packages/core/agent_framework/_agents.py | _merge_options |
| 名字合法化 | python/packages/core/agent_framework/_agents.py | _sanitize_agent_name |
| 单次 run 上下文类型 | python/packages/core/agent_framework/_agents.py | _RunContext |
| 流式响应收尾 | python/packages/core/agent_framework/_agents.py | RawAgent._parse_streaming_response |
| chat client 接口契约 | python/packages/core/agent_framework/_clients.py | SupportsChatGetResponse |
| chat client 抽象基类 | python/packages/core/agent_framework/_clients.py | BaseChatClient |
| 供应商必实现的方法 | python/packages/core/agent_framework/_clients.py | BaseChatClient._inner_get_response |
| client 转 agent 便捷法 | python/packages/core/agent_framework/_clients.py | BaseChatClient.as_agent |
| 能力探测协议族 | python/packages/core/agent_framework/_clients.py | SupportsCodeInterpreterTool 等 |
| 统一消息 | python/packages/core/agent_framework/_types.py | Message |
| 输入规整 | python/packages/core/agent_framework/_types.py | normalize_messages |
| chat / agent 响应 | python/packages/core/agent_framework/_types.py | ChatResponse / AgentResponse |
| 流式增量转换 | python/packages/core/agent_framework/_types.py | map_chat_to_agent_update |
| 流式包装器 | python/packages/core/agent_framework/_types.py | ResponseStream |
| 会话容器 | python/packages/core/agent_framework/_sessions.py | AgentSession |
| 上下文 provider 基类 | python/packages/core/agent_framework/_sessions.py | ContextProvider |
| 历史 provider 基类 | python/packages/core/agent_framework/_sessions.py | HistoryProvider |
| 内存历史(默认) | python/packages/core/agent_framework/_sessions.py | InMemoryHistoryProvider |
| 本地历史哨兵 id | python/packages/core/agent_framework/_sessions.py | LOCAL_HISTORY_CONVERSATION_ID |
| 每次调用持久化 middleware | python/packages/core/agent_framework/_sessions.py | PerServiceCallHistoryPersistingMiddleware |
接着读: 工具怎么定义、怎 么进循环、middleware 怎么拦截 → 第 02 章:工具、中间件与技能。整套库的全景与阅读地图 → 总览 index。