功能设计
my-agent-llm
是整个系统的「模型边界层」——负责抹平不同大模型厂商(OpenAI、DeepSeek、Anthropic
等)在网络通信、数据协议、流式输出、工具调用与思维链等方面的方言差异。
它让上层的
my-agent-core(核心调度微内核)永远只认识中立规范的接口与实体(如
Message、ToolCall、StreamEvent),完全无需感知背后是哪家供应商,也不用在业务微内核里承担底层的网络拼装脏活。
1 2 3 4 5 6 7 8 9 10 11 12 my-pi-agent 三层架构全景 │ ├── my-coding-agent 产品层:编码工具集 (read/write/edit/bash) + MCP 扩展 │ │ │ ▼ 依赖 ├── my-agent-core 框架层:无状态 ReAct 微内核 (loop.py) + 会话树 + Hook 门禁 │ │ │ ▼ 依赖 └── my-agent-llm ★ 模型边界层:统一门面 + 结构化 ToolCall + 统一累加器 (StreamAccumulator) + Provider 高阶事件流 │ ▼ 驱动 openai / anthropic / deepseek SDK / REST API
包结构
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 packages/my-agent-llm/ ├── pyproject.toml # name=my-agent-llm,src 布局 + hatchling ├── .env.example ├── src/my_agent_llm/ │ ├── __init__.py # 导出 LLM, Config, Message, Response, StreamChunk, StreamEvent 体系 │ ├── client.py # LLM 统一门面(chat / stream / achat / achat_stream / astream_events) │ ├── config.py # Config 不可变配置模型(pydantic frozen) │ ├── models.py # Message / Response / StreamChunk / ToolCall / TurnOutcome │ ├── events.py # [对标 Tau] 高阶流式事件契约 (StreamStart/TextDelta/ThinkingDelta/StreamDone/StreamError) │ ├── stream.py # [对标 Tau] 统一流式累加与规范化器 (StreamAccumulator) │ └── providers/ │ ├── _base.py # Provider 抽象契约 (ABC,包含 astream_events 默认实现) │ ├── registry.py # PROVIDER_REGISTRY 供应商注册表 │ ├── openai.py # OpenAIProvider (标准 OpenAI 协议转换) │ ├── deepseek.py # DeepSeekProvider (继承 OpenAI + reasoning_content 推理链提取) │ └── anthropic.py # AnthropicProvider (原生 dict 参数直传 + web_search 原生工具增强) └── tests/
供应商模块与分层设计
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 ┌─────────────────────────────────────────────────────────┐ │ LLM 统一门面(Facade) │ client.py │ chat / stream / achat / achat_stream / astream_events │ └────────────────────────────┬────────────────────────────┘ │ 按 config.provider 查表实例化 ┌────────────────────────────▼────────────────────────────┐ │ Provider 抽象契约(ABC) │ providers/_base.py └─────────────────────────────────────────────────────────┘ │ │ │ OpenAIProvider DeepSeekProvider AnthropicProvider (基准协议适配) (继承 OpenAI + 思考链) (原生 dict Block 互转) └──────────────┬──────────────────────────────┬──────┘ │ 汇入统一规范化累加流水线 ▼ ┌─────────────────────────────────────────────────────────┐ │ StreamAccumulator (流式累加与事件生成器) │ stream.py │ • 维护单一权威 partial: Message │ │ • 逐字累加文本与 Thinking 思维链 │ │ • 消化协作式取消 (CancellationToken) │ │ • 网络异常 Never-Throw 封装为 StreamErrorEvent │ │ • 终态产出合法完型的 StreamDoneEvent │ └─────────────────────────────────────────────────────────┘
基类契约与默认流式事件分发
_base.py 采用 ABC
作为父类,声明四大经典契约方法(chat、stream、achat、achat_stream):
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 class Provider (ABC ): """各 provider 的统一接口。""" @abstractmethod def __init__ (self, config: Config ): ... @abstractmethod def chat (self, messages: list [Message], *, model: str , tools: list [dict ] | None = None , **kwargs ) -> Response: ... @abstractmethod def stream (self, messages: list [Message], *, model: str , tools: list [dict ] | None = None , **kwargs ) -> Iterator[StreamChunk]: ... @abstractmethod async def achat (self, messages: list [Message], *, model: str , tools: list [dict ] | None = None , **kwargs ) -> Response: ... @abstractmethod async def achat_stream (self, messages: list [Message], *, model: str , tools: list [dict ] | None = None , **kwargs ) -> AsyncIterator[StreamChunk]: ...
[核心增强]
高阶事件流接口:astream_events
对标 Tau
(tau_ai/provider.py) ,基类提供了开箱即用的高阶流式事件流实现。子类只需实现标准的
achat_stream,基类自动通过 StreamAccumulator
将其升格为安全的高阶事件流:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 async def astream_events ( self, messages: list [Message], *, model: str , tools: list [dict ] | None = None , signal: Any | None = None , **kwargs, ) -> AsyncIterator[StreamEvent]: """异步高阶流式事件流(对标 Tau stream_response)。""" from ..stream import StreamAccumulator acc = StreamAccumulator() async for ev in acc.stream( self .achat_stream(messages, model=model, tools=tools, **kwargs), signal=signal, ): yield ev
核心数据模型与中立化抽象
(models.py)
为了消灭上层调度循环与底层各家 API 之间的强耦合,模型层在
models.py
中定义了五个核心数据模型。它们构成了整个框架通用的“世界语”,抹平一切供应商方言差异。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 ┌─────────────────────────────────────────────────────────────────────────────┐ │ models.py 核心数据实体与职责分工 │ ├─────────────────────────────────────────────────────────────────────────────┤ │ │ │ 1. TurnOutcome & normalize_finish_reason │ │ • 中立化的单轮终止原因枚举(抹平 "stop", "end_turn", "tool_use" 等方言)│ │ │ │ 2. Message │ │ • 框架通用的对话消息载体 (role + content + metadata) │ │ │ │ 3. ToolCall │ │ • 结构化工具调用第一公民 (id, name, args: dict, error) │ │ • 内置 @model_validator 自动解包与 Never-Throw 容错 │ │ • 彻底终结 4 重 JSON 编解码死循环 (4x JSON Ping-Pong) │ │ │ │ 4. Response │ │ • 单轮完整响应实体 (content, tool_calls, usage, reasoning, outcome) │ │ • 提供 to_message() 一键转为标准 Message,消灭手动拼装工人代码 │ │ │ │ 5. StreamChunk │ │ • 流式增量切片 (content 增量 + 末块携带完整的 response 实体) │ │ │ └─────────────────────────────────────────────────────────────────────────────┘
1.
中立化单轮终止状态机:TurnOutcome 与
normalize_finish_reason
结构与字段
1 2 3 4 5 6 7 8 9 class TurnOutcome (str , Enum): """Provider 中立的单轮终止原因(对标 pig-llm)。""" COMPLETED = "completed" TOOL_CALLS = "tool_calls" LENGTH = "length" CONTENT_FILTER = "content_filter" ABORTED = "aborted" PROVIDER_ERROR = "provider_error" UNKNOWN = "unknown"
作用与设计痛点
痛点 :各大模型厂商返回的 finish_reason
各说各话。OpenAI 正常说完叫 "stop",Anthropic 叫
"end_turn";调工具时 OpenAI 叫
"tool_calls",Anthropic 叫
"tool_use"。上层如果用字符串比对,极易漏判且充斥硬编码。
解法 :通过
normalize_finish_reason(reason: str | None, has_tool_calls: bool = False)
函数,统一经由预定义的 6 个状态集合进行归一化。即使厂商返回
None,只要存在工具调用也能自动修正推断为
TurnOutcome.TOOL_CALLS。上层只需通过
response.outcome == TurnOutcome.TOOL_CALLS
即可做强类型判断。
2.
统一对话消息载体:Message
结构与字段
1 2 3 4 5 class Message (BaseModel ): """统一消息:role + content + 附加元数据。""" role: Literal ["system" , "developer" , "user" , "assistant" , "tool" ] content: str metadata: dict [str , Any ] | None = None
作用与设计细节
结构说明 :
role:使用 Literal
限定标准角色,严格对齐工业级智能体角色定义;
content:消息纯文本正文;
metadata:承载随路扩展信息,包括:tool_calls(工具调用列表)、tool_call_id(Tool
消息对应的调用标识)、is_error(工具执行报错标记)、reasoning_content(思维链)、usage(Token
计费消耗)以及 stop_reason。
作用 :贯穿会话全流程的原子消息结构。不仅是内存
ReAct 循环流转的标准实体,也是 SessionEntry
持久化落盘的核心负载。
结构与字段
1 2 3 4 5 6 7 8 9 10 11 12 class ToolCall (BaseModel ): """统一结构化工具调用对象(对标 Tau / Pi)。""" id : str name: str args: dict [str , Any ] = Field(default_factory=dict ) error: str | None = None @model_validator(mode="before" ) @classmethod def _normalize_wire_dict (cls, data: Any ) -> Any : ... def to_wire_dict (self ) -> dict [str , Any ]: ...
核心设计与架构革新
原生 args: dict 交付(消灭 4x JSON
Ping-Pong) :
历史痛点 :过去模型层只给 JSON 字符串 ➔
核心层反序列化为字典进行 Hook 审查改参 ➔ 调注册表又序列化为字符串 ➔
注册表为了执行真实工具函数又反序列化为字典。整整 4 次
json.dumps 和 json.loads
往返折腾 !
革新成果 :参数在模型边界层内部反序列化完毕,交付给核心微内核的就是规整的原生
Python 字典。Anthropic 原生 SDK 返回的 block.input
字典被直接传递,微内核与注册表直接消费字典,全程 0 行多余 JSON
编解码 !
@model_validator(mode="before")
自动解包与容错 :
自适应兼容 :如果输入的是 OpenAI 原始嵌套传输格式
{"id": "...", "function": {"name": "...", "arguments": "..."}},前置验证器自动解包为平铺的
id, name, args 结构;
Never-Throw 崩溃防御 :当大模型输出畸形、非法 JSON
字符串时,绝不上抛异常打崩程序 ,而是将
args 兜底置为空字典,并将解析异常信息记录在
ToolCall.error 字段中。核心调度微内核据此转化为标准的
ToolResult(ok=False, error=...),让大模型在下一轮自我修复!
to_wire_dict()
协议导出 :提供反向导出为 OpenAI Wire
形状的能力,兼顾了与外部系统和旧协议的无缝对接。
4.
单轮完整响应实体:Response
结构与字段
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 class Response (BaseModel ): """统一响应:文本 + 工具调用 + usage + reasoning。""" content: str model: str tool_calls: list [ToolCall] | None = None reasoning_content: str | None = None usage: dict [str , int ] | None = None finish_reason: str | None = None @property def outcome (self ) -> TurnOutcome: return normalize_finish_reason(self .finish_reason, bool (self .tool_calls)) def to_message ( self, role: Literal ["system" , "developer" , "user" , "assistant" , "tool" ] = "assistant" , stop_reason: str | None = None , ) -> Message: ...
作用与设计细节
作用 :用于非流式调用(chat/achat),以及流式末块的完整交付。
outcome 动态属性 :内部自动调用
normalize_finish_reason,对外暴露统一的
TurnOutcome 枚举,供状态机做跳转判断。
to_message()
一键转换 :模型层自己最清楚哪些属性应该映射到
Message.metadata。通过一键转换,彻底消灭调度微内核手工拼装
metadata 字典的工人代码。
5.
流式增量切片:StreamChunk
结构与字段
1 2 3 4 5 6 7 8 9 10 11 class StreamChunk (BaseModel ): """流式增量块:文本增量 + 末块携带完整 tool_calls 与终态已拼装好的 Response。""" content: str finish_reason: str | None = None tool_calls: list [ToolCall] | None = None usage: dict [str , int ] | None = None metadata: dict [str , Any ] | None = None response: Response | None = None @property def outcome (self ) -> TurnOutcome | None : ...
作用与双轨设计
轻量流式打字 :流式传输过程中,各块携带增量文本(content),供打字机或终端实时吐字渲染;
末块终态交付 :在流式最后一个分片上,携带组装好的完整
response: Response 实体。下游的
StreamAccumulator
或微内核无需手写字符串累加,直接读取末块的 response
即可获得完整结果。
模型层流式调用的底层实现与三级异步接力机制
很多初学者容易觉得大模型流式调用抽象,主要是因为没有看清网络物理层到
Python 运行栈的接力过程 。
流式调用的物理本质非常朴素: >
大模型服务端每吐出一个字,就通过 Python 的
async for 和 yield
异步生成器像“接力棒”一样,毫秒不差地层层往上传递到终端屏幕上。
1. 物理本质:单次阻塞 vs
流式管道 (SSE)
1 2 3 4 5 6 7 8 9 10 11 【普通非流式调用 (chat / achat)】 你发请求 ──► OpenAI 计算 3 秒(程序挂起等待) ──► 一次性返回整段大 JSON 字符串 (缺点:用户要对着空白屏幕干等好几秒,无法实现打字机效果) 【流式调用 (stream / achat_stream)】 你发请求 (带参数 stream=True) ──► OpenAI 保持 HTTP 长连接不断开 (Server-Sent Events: SSE),像水管流水一样,每算出一个 Token 就推一个微小分块: ├── 第 1 个分块: "我" ├── 第 2 个分块: "正在" ├── 第 3 个分块: "帮您" └── 最后一个分块: [DONE] (断开连接)
2.
流式的“三级接力赛”(从物理网卡到屏幕打字机)
整个流式调用由 4 个关键模块
紧密串联,形成了一条无阻塞的异步直通管道:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 ┌─────────────────────────────────────────────────────────────────────────┐ │ 【第 1 棒:物理网卡 ➔ Provider】(openai.py / anthropic.py) │ │ 负责连接官方 SDK 的底层 SSE 网络流,逐字 yield StreamChunk │ └────────────────────────────────────┬────────────────────────────────────┘ │ yield StreamChunk(content="我") ▼ ┌─────────────────────────────────────────────────────────────────────────┐ │ 【第 2 棒:Provider ➔ 累加中枢】(stream.py: StreamAccumulator) │ │ 在模型边界内维持单一事实快照,把原始 Chunk 升格为 StreamEvent │ └────────────────────────────────────┬────────────────────────────────────┘ │ yield TextDeltaEvent(delta="我", partial=...) ▼ ┌─────────────────────────────────────────────────────────────────────────┐ │ 【第 3 棒:累加中枢 ➔ 调度微内核】(my-agent-core/loop.py: _assistant_turn)│ │ 将模型层的高阶事件,转译为 Agent 框架的生命周期事件 MessageUpdate │ └────────────────────────────────────┬────────────────────────────────────┘ │ yield MessageUpdate(chunk=...) ▼ ┌─────────────────────────────────────────────────────────────────────────┐ │ 【第 4 棒:微内核 ➔ 终端打字机】(main.py / UI 界面) │ │ sys.stdout.write(chunk.content) 并在屏幕上实时打印出这个字! │ └─────────────────────────────────────────────────────────────────────────┘
3. 逐级源码追踪:一个字
"我" 是如何穿透系统的?
第一棒:在
providers/openai.py 里(发起真实 SSE 连接)
1 2 3 4 5 6 7 8 9 10 11 12 13 stream = await self .async_client.chat.completions.create( model="gpt-4o" , messages=..., stream=True , ) async for chunk in stream: delta = chunk.choices[0 ].delta if delta.content: yield StreamChunk(content=delta.content)
第二棒:在
stream.py 里(累加器维持状态,升格为事件)
1 2 3 4 5 6 7 8 acc = StreamAccumulator() async for chunk in source: for ev in self .feed(chunk): yield ev
第三棒:在
my-agent-core/loop.py 里(调度微内核转译)
1 2 3 4 5 async for ev in llm.astream_events(...): if isinstance (ev, TextDeltaEvent): yield MessageUpdate(message=ev.partial, chunk=StreamChunk(content=ev.delta))
第四棒:在终端入口
main.py 里(打字机实时吐字)
1 2 3 4 5 6 7 async for event in agent.prompt_stream("你好" ): if isinstance (event, MessageUpdate): if event.chunk and event.chunk.content: sys.stdout.write(event.chunk.content) sys.stdout.flush()
4.
为什么它不会卡死?(异步生成器的穿透力)
这就是 Python async for 配合
yield(异步生成器) 的魔力所在: 1.
普通阻塞函数 :必须把整层循环跑完,攒成一个大列表
return list,下游才能拿到数据; 2.
异步生成器(yield) : - 只要远端 OpenAI
吐了一个字,底层的 yield
就会像多米诺骨牌一样瞬间触发上层所有的
yield ; - 这个字在 1 毫秒内
就穿透了四层调用栈,被打字机直接刷在屏幕上; -
打印完后,整个程序立刻让出 CPU 线程(进入
await),等待网卡下一个数据包的到来。
5.
核心代码透视:为什么是
for ev in self.feed(chunk): yield ev?
在阅读 StreamAccumulator.stream
的源码时,很多人会困惑这行代码: 1 2 3 async for chunk in source: for ev in self .feed(chunk): yield ev
既然 chunk
是一个个进来的,为什么还要对 self.feed(chunk)
的结果再套一层 for 循环?
答案在于:底层网络分片(StreamChunk)与高阶领域事件(StreamEvent)之间,并不是绝对的一对一,而是“一对多”的动态映射关系!
(1) 1 对 N
映射:单个 Chunk 可能拆出 0 个、1 个或多个事件
场景 A(首包到达:1 变
2) :当大模型刚连上,吐出第一个字 "你"
时,状态机既要发出“流式启动”广播(StreamStartEvent),又要发出“正文打字”增量(TextDeltaEvent),单个
Chunk 必须同时产出 2 个事件 ;
场景 B(混合分块:1 变 2) :当一个 Chunk
同时携带了思考链和正文(例如 reasoning_content="分析完毕"
且 content="答案是"),feed 会顺序产出
ThinkingDeltaEvent 和 TextDeltaEvent;
场景 C(多工具并发:1 变
N) :若模型单轮输出了多个并行工具调用,feed
会遍历产出多个独立的 ToolCallDoneEvent;
场景 D(纯统计或心跳包:1 变 0) :流式结束时的纯
usage 统计包没有文本增量,feed 返回空列表
[],循环走 0 次,向外界过滤无意义的空事件噪音。
(2)
架构解耦:“同步纯计算”与“异步 I/O”分离
feed(chunk)
是同步纯函数 :专职负责内存计算(字符串累加、JSON
解析、状态机判定),耗时 < 1
微秒,零网络阻塞。这使得单测无需启动复杂的事件循环,直接
events = acc.feed(chunk) 即可瞬间断言!
stream(source)
是异步搬运工 :专职处理网络等待(async for chunk in source)与协作取消监听(is_cancelled())。
(3)
语法机制:异步生成器中的“拍平展开(Flatten)”
在 Python 的同步函数中,我们可以写 yield from [a, b]
把列表元素逐个送出; 但在 Python
的异步生成器(async def) 中,语言语法并不支持
yield from ! 因此: 1 2 for ev in self .feed(chunk): yield ev
正是异步生成器中标准、正统的“扁平化拍平(Flatten)”写法 。
生动比喻 : -
底层网卡送来的是“箱子”(chunk) ; -
self.feed(chunk)
是“开箱拆包” (同步动作:从箱子里掏出 0 个、1
个或多个“标准零件”); -
for ev in self.feed(chunk): yield ev
就是“把拆出来的零件一个个排上传送带” ,推给终端屏幕!
对标
Tau 的重大架构升级:StreamEvent 体系与 StreamAccumulator
在阶段 18~21 的演进中,我们深度对标了 Tau
(tau_agent/loop.py 与
tau_ai/stream.py) ,彻底消除了微内核
loop.py 中的职责撕裂。
为什么低阶
StreamChunk 会导致微内核职责污染?
在早期设计中,模型层只抛出原始的
StreamChunk,导致调度微内核 loop.py 的
_assistant_turn 被迫充当“包工头”: 1. 自己维护
started = False 哨兵来判定何时向 UI 发送
MessageStart; 2. 自己手写
fallback_content += chunk.content 累加字符串; 3. 自己手写
try...except 捕获断网异常并拼装错误文本; 4.
退出时自己判断有没有终态实体,手写
Message(role="assistant", ...)。
主调度循环充斥着大模型网络与字符串拼装的脏活,代码膨胀至 120+
行!
解决方案:模型层自闭环(StreamAccumulator)
对标 Tau 的 canonicalize_provider_stream,我们在
my_agent_llm.stream 实现了通用累加器:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 原始网络 Chunks (来自各大 Provider) │ ▼ ┌─────────────────────────────────────────────────────────────┐ │ StreamAccumulator (模型层内部的装配与质检车间) │ │ │ │ 1. [首块通知]: 首个 chunk 到达 ➔ yield StreamStartEvent │ │ 2. [打字增量]: chunk.content ➔ 累加并 yield TextDeltaEvent │ │ 3. [思维推理]: reasoning ➔ 累加并 yield ThinkingDeltaEvent │ │ 4. [工具就绪]: tool_calls ➔ 组装并 yield ToolCallDoneEvent │ │ 5. [异常隔离]: 监听 signal 与网络异常,Never-Throw 封装为 │ │ yield StreamErrorEvent(error=...) │ │ 6. [终态交付]: yield StreamDoneEvent(message=..., usage=...)│ └─────────────────────────────────────────────────────────────┘ │ ▼ 产出高阶流式事件流 (StreamEvent) ┌─────────────────────────────────────────────────────────────┐ │ 框架调度微内核 (_assistant_turn 仅约 25 行纯转译逻辑!) │ └─────────────────────────────────────────────────────────────┘
微内核
_assistant_turn 的断崖式精简
重构后,my-agent-core/src/my_agent_core/loop.py 中的
_assistant_turn 实现了极致的优雅与纯粹:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 async def _assistant_turn ( *, llm: Any , view: list [Message], tool_schemas: list [dict [str , Any ]], model: str | None = None , signal: CancellationToken | None = None , context_manager: Any | None = None , ) -> AsyncIterator[Event]: """专职大模型推理车间:纯粹转译模型层产出的高阶 StreamEvent(对标 Tau _assistant_events)。""" async for ev in llm.astream_events( messages=view, tools=tool_schemas, model=model, signal=signal ): if isinstance (ev, StreamStartEvent): yield MessageStart(ev.partial) elif isinstance (ev, TextDeltaEvent): yield MessageUpdate(message=ev.partial, chunk=StreamChunk(content=ev.delta)) elif isinstance (ev, ThinkingDeltaEvent): yield MessageUpdate( message=ev.partial, chunk=StreamChunk(content="" , metadata={"reasoning_content" : ev.delta}), ) elif isinstance (ev, ToolCallDoneEvent): yield MessageUpdate( message=ev.partial, chunk=StreamChunk(content="" , tool_calls=[ev.tool_call]), ) elif isinstance (ev, StreamDoneEvent): if ev.usage and context_manager is not None : with contextlib.suppress(Exception): context_manager.record_usage(ev.usage) yield MessageEnd(ev.message) elif isinstance (ev, StreamErrorEvent): yield MessageEnd(ev.error)
对比成效 : - 0 行
started 哨兵变量判断; - 0 行
fallback_content += chunk.content 手工字符串拼接; -
0 行 try...except 网络异常拼装; -
0 行 手动 to_message 降级重构; -
真正的职责分离 :模型层全权负责模型通信、Token
累加、容错降级与状态快照;框架微内核仅负责业务状态机流转与领域事件广播!
模型层高阶流式事件清单与端到端数据流转全景
在完成 Tau
对标后,模型层对外暴露的不再是零碎不可靠的原始网络分片,而是一套定义在
my_agent_llm.events 中的高阶事件契约流。
一、模型层流式事件清单
(events.py)
所有事件继承自 StreamEvent
基类(dataclass(frozen=True)),保证在异步管道流动中的不可变性:
1 2 3 4 5 6 7 8 9 10 11 StreamEvent (高阶流式事件基类) │ ┌──────────────────┼──────────────────┬──────────────────┐ ▼ ▼ ▼ ▼ StreamStartEvent TextDeltaEvent ThinkingDeltaEvent ToolCallDoneEvent (首字启动) (正文文本增量) (思维链增量) (工具参数完型) │ ┌─────────────┴─────────────┐ ▼ ▼ StreamDoneEvent StreamErrorEvent (正常完型交付) (异常/取消安全交付)
事件类名
携带核心属性
触发时机与语义
设计目的与收益
StreamStartEvent
partial: Message
网卡读出首个 Token 或连接建立时触发
通知下游立即初始化对话气泡,挂起转圈动画/光标。消除首字延迟期间前端的无响应感。
TextDeltaEvent
delta: strpartial: Message
接收到模型吐出的正文 Token 增量
驱动终端打字机逐字输出;随路携带最新的全量累积
partial,供下游无需自己开缓冲区。
ThinkingDeltaEvent
delta: strpartial: Message
接收到 DeepSeek-R1 / Claude 3.7
的思考链分片
将推理思维过程作为一等公民,支持前端独立折叠框展示,与最终正文完全解耦。
ToolCallDeltaEvent
index: intdelta: strpartial: Message
工具调用参数分片流式拼接过程
提供更细粒度的工具参数流式可观测性。
ToolCallDoneEvent
index: inttool_call: ToolCallpartial: Message
单个工具调用的 JSON
参数解析并验证完毕
提前通知调度层工具已完整就绪,使 UI
能在流式未完全结束前提前渲染“准备运行”态。
StreamDoneEvent
message: Messageusage: dict[str, int]
底层网络流正常完型
交付合法定型的完整
Message,并携带真实 Token
usage,为上下文压缩提供权威校准锚点。
StreamErrorEvent
error: Messagestop_reason: strexc: Exception
网络中断、超时、429
或主动取消(signal)
Never-Throw
保证 :异常不崩溃,转化为带
stop_reason="error"/"cancelled"
的合法消息,驱动后续自愈。
二、端到端数据流转全景(从用户输入到工具执行)
我们以一次真实的“模型先思考 ➔ 吐正文 ➔
发起工具调用” 为例,展示事件与数据在 应用层 ➔ 框架微内核
➔ 模型边界层 ➔ 底层物理网络 ➔ 回流调度
中的完整自上而下运转全景:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 1: 用户发起交互与轮次初始化 │ │ ├─ [应用层 / Agent] 用户提问:"帮我分析代码并查询北京天气" │ │ ├─ [微内核 run_agent_loop] 生成 user Message 并准备当前轮次 view 视图 │ │ └─ [微内核 _assistant_turn] 调用: async for ev in llm.astream_events(...) │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 2: 建立网络连接,读到首包 (StreamStartEvent) │ │ ├─ [底层网卡] 收到首个 SSE Chunk │ │ ├─ [模型层 StreamAccumulator] 识别 started=False │ │ │ └─► 触发 yield StreamStartEvent(partial) │ │ ├─ [微内核 _assistant_turn] 接收并转译 │ │ │ └─► 转译 yield MessageStart(assistant) │ │ └─ [终端 UI / 前端] │ │ └─► 立即弹出空白助手消息气泡,挂起转圈动画与光标 │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 3: 吐出思维链碎片 (ThinkingDeltaEvent,如 DeepSeek-R1 / Claude 3.7) │ │ ├─ [底层网卡] 吐出 reasoning_content: "正在分析用户的查询意图..." │ │ ├─ [模型层 StreamAccumulator] 累加至 metadata["reasoning_content"] │ │ │ └─► 触发 yield ThinkingDeltaEvent(delta, partial) │ │ ├─ [微内核 _assistant_turn] 接收并转译 │ │ │ └─► 转译 yield MessageUpdate(message=partial) │ │ └─ [终端 UI / 前端] │ │ └─► 前端在独立折叠框中实时打字滚动显示思考过程 │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 4: 吐出正文打字机增量 (TextDeltaEvent) │ │ ├─ [底层网卡] 吐出 content: "我来帮您查询北京今天的天气:" │ │ ├─ [模型层 StreamAccumulator] 累加至 self.content │ │ │ └─► 触发 yield TextDeltaEvent(delta, partial) │ │ ├─ [微内核 _assistant_turn] 接收并转译 │ │ │ └─► 转译 yield MessageUpdate(message=partial) │ │ └─ [终端 UI / 前端] │ │ └─► 打字机丝滑逐字打出正文 │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 5: 工具参数拼接与完型就绪 (ToolCallDoneEvent) │ │ ├─ [底层网卡] 吐出 tool_calls 碎片 (name, arguments JSON) │ │ ├─ [模型层 StreamAccumulator] 完成 JSON 反序列化为原生 Python 字典 │ │ │ └─► 触发 yield ToolCallDoneEvent(tool_call, partial) │ │ ├─ [微内核 _assistant_turn] 接收并转译 │ │ │ └─► 转译 yield MessageUpdate(tool_calls=...) │ │ └─ [终端 UI / 前端] │ │ └─► 界面高亮提示:“工具 get_weather 即将调用...” │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 6: 网络流正常完型与定型交付 (StreamDoneEvent) │ │ ├─ [底层网卡] 流式结束,附带官方真实 Usage 开销 │ │ ├─ [模型层 StreamAccumulator] 一键 to_message() 封装终态实体 │ │ │ └─► 触发 yield StreamDoneEvent(message, usage) │ │ ├─ [微内核 _assistant_turn] 接收并终结 │ │ │ ├─► context_manager.record_usage(usage) (锚定真实 Token 计费) │ │ │ └─► 转译 yield MessageEnd(message) │ │ └─ [终端 UI / 前端] │ │ └─► 关闭打字机光标,固化 Markdown 渲染结果与代码高亮 │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 7: 进入工具批处理执行流水线 (_execute_tools_turn) │ │ ├─ 拿到合法的完整 assistant: Message 实体 │ │ ├─ [Preflight 阶段] 率先广播 ToolExecutionStart (供 UI 渲染执行态) │ │ ├─ [审批阶段] before_tool_call (ToolCallHook) 安全改参或审批放行 │ │ ├─ [执行阶段] ToolRegistry.execute_batch 并发批执行原生字典工具 │ │ ├─ [改写阶段] after_tool_call (ToolResultHook) 篡改出参或捕获异常 │ │ └─ [闭环阶段] 广播 ToolExecutionEnd ➔ 产出 role="tool" 消息 ➔ 闭环 TurnEnd │ └─────────────────────────────────────────────────────────────────────────────┘
三、意外与中断自愈流转(例如用户中途按
Ctrl+C 取消)
如果在流式生成的任何一刻发生意外中断(如用户按
Ctrl+C、或网络超时断开),系统通过显式取消信号(CancellationToken)实现自上而下的绝对优雅闭环:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 1: 外部触发取消 │ │ └─► 外部调用 agent.abort() 或捕获 SIGINT ➔ CancellationToken.cancel() = True │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 2: 模型层 StreamAccumulator 协作感知与安全打断 │ │ ├─ Accumulator 在下一个 Chunk 循环中检测到 is_cancelled() == True │ │ ├─ 立即安全 break 跳出底层 HTTP 连接接收,不挂起僵尸连接 │ │ ├─ 保留打断前已经生成出来的半截残句 (partial.content) │ │ └─ 给消息打上元数据标签 metadata={"stop_reason": "cancelled"} │ │ (若为底层网络抛异常,则标记为 "error" 并附带异常信息) │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 3: 模型层向微内核发射 StreamErrorEvent │ │ └─► yield StreamErrorEvent(error=error_msg, stop_reason="cancelled") │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 4: 微内核 _assistant_turn 接收并闭环消息 │ │ └─► 优雅转译 yield MessageEnd(error_msg),绝不向上抛出未捕获异常打崩主程序! │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 步骤 5: 调度微内核 run_agent_loop 捕获并启动转录本自愈 │ │ ├─ 发现 assistant 消息携带 stop_reason="cancelled" │ │ ├─ 检查到该消息带有尚未执行的 tool_calls │ │ ├─ [核心自愈]: 触发 _synthesize_interrupted_tool_calls │ │ │ └─► 100% 自动为悬空调用合成 role="tool", │ │ │ content="Tool call interrupted by user", is_error=True │ │ └─ 写入 Session 树持久化落盘,严格闭环发射配对的 TurnEnd │ └──────────────────────────────────────┬──────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ 终态结果: 系统安全退出,且会话拓扑时刻合法! │ │ └─► 彻底杜绝下一次用户提问时触发大模型 API 400 校验死锁! │ └─────────────────────────────────────────────────────────────────────────────┘
LLM 门面类 (client.py)
在吸纳了 Tau 的无多余包装思想后,LLM
门面类被重构为纯粹、对称且强类型的统一入口:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 class LLM : """统一 LLM 客户端门面:严格接收强类型 Config 配置,多态路由至具体 Provider。""" def __init__ (self, config: Config ) -> None : """标准工程构造:严格接收 Config 实例,彻底废除 kwargs 散参构造。""" if not isinstance (config, Config): raise TypeError(f"LLM expects a Config instance, got {type (config).__name__} " ) if config.provider not in PROVIDER_REGISTRY: raise ValueError(f"Unknown provider '{config.provider} '..." ) if not config.api_key: raise ValueError(f"No API key for provider: {config.provider} " ) provider_cls = PROVIDER_REGISTRY[config.provider] self ._provider: Provider = provider_cls(config) self .config = config @property def model (self ) -> str : ... def _resolve_kwargs (self, kwargs: dict [str , Any ] ) -> dict [str , Any ]: """集中收敛参数级联回退:调用传参 > Config 默认值,彻底消灭分散的 if-else。""" opts = dict (kwargs) if self .config.temperature is not None : opts.setdefault("temperature" , self .config.temperature) if self .config.max_tokens is not None : opts.setdefault("max_tokens" , self .config.max_tokens) return opts def chat (self, messages: list [Message], *, tools: list [dict ] | None = None , model: str | None = None , **kwargs: Any ) -> Response: ... def stream (self, messages: list [Message], *, tools: list [dict ] | None = None , model: str | None = None , **kwargs: Any ) -> Iterator[StreamChunk]: ... async def achat (self, messages: list [Message], *, tools: list [dict ] | None = None , model: str | None = None , **kwargs: Any ) -> Response: ... async def achat_stream (self, messages: list [Message], *, tools: list [dict ] | None = None , model: str | None = None , **kwargs: Any ) -> AsyncIterator[StreamChunk]: ... async def astream_events ( self, messages: list [Message], *, tools: list [dict ] | None = None , model: str | None = None , signal: Any | None = None , **kwargs: Any , ) -> AsyncIterator[StreamEvent]: """异步高阶流式事件流:直接产出 StreamStart/TextDelta/ThinkingDelta/StreamDone/StreamError。""" opts = self ._resolve_kwargs(kwargs) async for ev in self ._provider.astream_events( messages, model=model or self .model, tools=tools, signal=signal, **opts ): yield ev
深度剖析工业级标杆:Tau
官方模型层的架构设计与实现精髓
在 Python 智能体生态中,Tau (tau-ai /
tau_agent)
的模型边界层被公认为架构设计最优雅、工业级成熟度最高的范本之一。
通过研读 Tau
的官方源码实现,我们可以系统性地总结出其四个极具启发性的核心设计:
1.
分层拓扑与包边界解耦(tau_ai vs
tau_agent)
Tau 在物理结构上将模型能力和智能体运行时做了绝对的物理拆分: -
tau_ai(模型通信与协议翻译 SDK) :专职处理
HTTP 建连、SSE 解析、各家模型私有协议转换(OpenAI Responses API /
ChatCompletions API、Anthropic Blocks、Google REST
等)、指数退避重试(Jitter)、Token 计费与首包延迟(TTFT)统计; -
tau_agent(便携运行时) :完全感知不到 HTTP
和底层
SDK,只对接供应商中立的高阶事件流(AssistantMessageEvent)。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 ┌────────────────────────────────────────────────────────────────────────┐ │ 【底层适配器】tau_ai/openai_compatible.py / anthropic.py / google.py │ │ • 建立 HTTP / SSE 网络连接,处理各家私有协议 │ │ • 内部发射原始轻量内部事件:ProviderTextDelta, ProviderToolCall 等 │ └───────────────────────────────────┬────────────────────────────────────┘ │ 原始 ProviderEvent 流 ▼ ┌────────────────────────────────────────────────────────────────────────┐ │ 【中央规范化流水线】tau_ai/stream.py: canonicalize_provider_stream │ │ • 维护单一事实来源:partial: AssistantMessage (有序块级内容) │ │ • 自动管理多通道(文本通道、思维通道、工具通道) │ │ • 拦截网络异常与取消信号,Never-Throw 封装为 AssistantErrorEvent │ │ • 终态自动组装 Usage、耗时统计(TTFT)与合法 Message │ └───────────────────────────────────┬────────────────────────────────────┘ │ 交付供应商中立的高阶事件流 ▼ ┌────────────────────────────────────────────────────────────────────────┐ │ 【上层调度微内核】tau_agent/loop.py: _assistant_events │ │ • 极其干净的 30 行转译逻辑: │ │ StartEvent ➔ MessageStart, DoneEvent ➔ MessageEnd, Delta ➔ Update │ │ • 核心循环零累加、零拼装、零异常处理! │ └────────────────────────────────────────────────────────────────────────┘
2. Tau
模型层的四大核心精妙设计
(1)
双层事件管道机制(Private ProviderEvent ➔ Public Canonical Event)
痛点 :底层各大供应商(OpenAI、Anthropic、Google、DeepSeek)的流式返回千奇百怪,有的甚至存在多个并发通道(例如
OpenAI 同时吐思维链和文本,Anthropic 带有 content_block 索引)。
Tau 解法 :
内层事件(_provider_events.py) :适配器只需发射简单的
ProviderTextDeltaEvent、ProviderThinkingDeltaEvent、ProviderToolCallEvent;
外层事件(provider_events.py) :由中央状态机
canonicalize_provider_stream
统一转译为高阶公共事件(AssistantStartEvent,
TextDeltaEvent, ThinkingDeltaEvent,
AssistantDoneEvent,
AssistantErrorEvent)。
收益 :接入新模型时,适配器只需要写一个几十行的最简解析器,剩下的重试、打字机事件、取消与状态累加全部由中央流水线自动赋能。
(2) 随路快照投影(Snapshot
in Delta Events)
在 Tau 中,所有增量事件都同时携带了累加快照: 1 2 3 class TextDeltaEvent (WireModel ): delta: str partial: AssistantMessage
-
收益 :下游观察者零计算负担 。如果终端打字机想打印,只看
delta;如果 Web
气泡或安全审核插件需要实时检查完整句子,直接读取
partial.text,绝不需要自己在下游开缓冲区累加字符串,实现单一数据源(Single
Source of Truth) 。
(3)
有序块级内容数据模型(Ordered Content Blocks)
1 2 class AssistantMessage (WireModel ): content: list [TextContent | ThinkingContent | ToolCall]
收益 :真实前沿模型(DeepSeek-R1、Claude
3.7、GPT-5)的行为是交替的:先思考 ➔ 吐一部分字 ➔ 调工具 ➔
拿到结果再思考 ➔ 输出最终答案 。纯文本 content: str
无法表达这种时间因果序列,而块级列表原生还原了模型的认知流,将 Thinking
与 ToolCall 升格为一等公民。
(4)
模型层自闭环容错与“错误即数据”(Never-Throw & Failures as
Data)
设计 :网络抖动、限流(429)、网关超时(504)由底层指数退避重试;重试耗尽或中途取消时,模型层绝不上抛崩溃异常 ,而是构建一条合法的
AssistantMessage(stop_reason="error" 或 "aborted", error_message="..."),包装在
AssistantErrorEvent 中传递给调度层。
收益 :异常不打崩程序,而是转化为结构化数据,供下游转录本自愈模块(tool_history.py)识别并自动补齐断头调用,彻底终结
API 400 校验死锁。
3. 架构收益全景对比表
核心维度
传统 Agent 框架(低阶 Chunk 模式)
Tau 架构(高阶事件流与累加自闭环)
微内核体量
调度微内核充当包工头,处理首字、累加字符串、捕获异常,膨胀到
120~200 行 意大利面条代码。
微内核 _assistant_events 仅
20~30 行 ,只有纯粹的 if isinstance
事件转译,一眼见底。
状态一致性
微内核在累加、UI
在累加、中间件在暂存,极易发生状态不一致或时序颠倒。
单一权威来源 ,模型层是唯一的装配车间,随路广播
partial 快照。
思维链支持
只能用各种 metadata
临时补丁到处拼接,展示与存储逻辑别扭。
原生支持 ThinkingDeltaEvent
与有序块,思维链是一等公民。
新模型接入
每个 Provider
都要考虑如何累加、如何处理异常与中断。
适配器只需抛出最简内部
delta,中央累加器自动赋能全量事件。
系统稳定性
网络波动或取消信号直接引发未捕获异常击穿
ReAct 循环。
Never-Throw
保证 ,失败与取消下沉为带规范 stop_reason
的实体,驱动后续自愈。
核心设计原则与架构不变式总结
Never-Throw Guarantee(模型层崩溃免疫) :
网络断开、超时或大模型产出畸形格式时,模型层绝不抛出未捕获异常打崩主程序,统一就地封装为带
stop_reason="error" 的
StreamErrorEvent,引导上层自愈或安全闭环;
Strict Timing Order Invariant(严密生命周期时序) :
流式事件严格遵循
StreamStartEvent ➔ *DeltaEvent... ➔ StreamDoneEvent / StreamErrorEvent
的不可逆确定性时序,彻底消灭首字延迟前提前发射或结束消息倒置的时序
Bug;
Structured First-Class
Citizens(结构化实体一等公民) :
参数在模型层完成反序列化,彻底消除 4 重 JSON 编解码死循环(4x JSON
Ping-Pong),直连微内核调度与工具执行。