my-pi-agent--模型层

功能设计

my-agent-llm 是整个系统的「模型边界层」——负责抹平不同大模型厂商(OpenAI、DeepSeek、Anthropic 等)在网络通信、数据协议、流式输出、工具调用与思维链等方面的方言差异。

它让上层的 my-agent-core(核心调度微内核)永远只认识中立规范的接口与实体(如 MessageToolCallStreamEvent),完全无需感知背后是哪家供应商,也不用在业务微内核里承担底层的网络拼装脏活。

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 作为父类,声明四大经典契约方法(chatstreamachatachat_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. 中立化单轮终止状态机:TurnOutcomenormalize_finish_reason

结构与字段

1
2
3
4
5
6
7
8
9
class TurnOutcome(str, Enum):
"""Provider 中立的单轮终止原因(对标 pig-llm)。"""
COMPLETED = "completed" # 正常生成结束 ("stop", "end_turn", "natural")
TOOL_CALLS = "tool_calls" # 发起工具调用 ("tool_calls", "tool_use", "function_call")
LENGTH = "length" # 超出最大 Token 限制 ("length", "max_tokens")
CONTENT_FILTER = "content_filter"# 命中内容审查拦截 ("safety", "blocked", "recitation")
ABORTED = "aborted" # 用户主动取消/中断 ("aborted", "cancelled")
PROVIDER_ERROR = "provider_error"# 厂商服务报错 ("error", "failed")
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 持久化落盘的核心负载。

3. 结构化工具调用实体:ToolCall

结构与字段

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]: ...

核心设计与架构革新

  1. 原生 args: dict 交付(消灭 4x JSON Ping-Pong)
    • 历史痛点:过去模型层只给 JSON 字符串 ➔ 核心层反序列化为字典进行 Hook 审查改参 ➔ 调注册表又序列化为字符串 ➔ 注册表为了执行真实工具函数又反序列化为字典。整整 4 次 json.dumpsjson.loads 往返折腾
    • 革新成果:参数在模型边界层内部反序列化完毕,交付给核心微内核的就是规整的原生 Python 字典。Anthropic 原生 SDK 返回的 block.input 字典被直接传递,微内核与注册表直接消费字典,全程 0 行多余 JSON 编解码
  2. @model_validator(mode="before") 自动解包与容错
    • 自适应兼容:如果输入的是 OpenAI 原始嵌套传输格式 {"id": "...", "function": {"name": "...", "arguments": "..."}},前置验证器自动解包为平铺的 id, name, args 结构;
    • Never-Throw 崩溃防御:当大模型输出畸形、非法 JSON 字符串时,绝不上抛异常打崩程序,而是将 args 兜底置为空字典,并将解析异常信息记录在 ToolCall.error 字段中。核心调度微内核据此转化为标准的 ToolResult(ok=False, error=...),让大模型在下一轮自我修复!
  3. 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 foryield 异步生成器像“接力棒”一样,毫秒不差地层层往上传递到终端屏幕上。

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
# 1. 调官方异步 SDK,传 stream=True
stream = await self.async_client.chat.completions.create(
model="gpt-4o",
messages=...,
stream=True, # ← 关键!要求服务端以 SSE 流式推送
)

# 2. 用 async for 从网卡套接字中逐个读出碎片
async for chunk in stream:
delta = chunk.choices[0].delta
if delta.content:
# 每从网络读到一个字,立刻 yield 出去!
yield StreamChunk(content=delta.content)

第二棒:在 stream.py 里(累加器维持状态,升格为事件)

1
2
3
4
5
6
7
8
acc = StreamAccumulator()

# 这里消费第一棒抛出来的原始流
async for chunk in source:
# 1. 把这个字累加到内部的 self.content 字符串里(蓄水)
# 2. 打包成带最新全量快照的高阶事件抛出去
for ev in self.feed(chunk):
yield ev # 抛出 TextDeltaEvent(delta="我", partial=Message(...))

第三棒:在 my-agent-core/loop.py 里(调度微内核转译)

1
2
3
4
5
# 这里消费第二棒抛出来的高阶事件
async for ev in llm.astream_events(...):
if isinstance(ev, TextDeltaEvent):
# 转译成框架通用的 MessageUpdate 事件,继续往外推!
yield MessageUpdate(message=ev.partial, chunk=StreamChunk(content=ev.delta))

第四棒:在终端入口 main.py 里(打字机实时吐字)

1
2
3
4
5
6
7
# 外部消费 agent.prompt_stream
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 会顺序产出 ThinkingDeltaEventTextDeltaEvent
  • 场景 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.pytau_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: str
partial: Message
接收到模型吐出的正文 Token 增量 驱动终端打字机逐字输出;随路携带最新的全量累积 partial,供下游无需自己开缓冲区。
ThinkingDeltaEvent delta: str
partial: Message
接收到 DeepSeek-R1 / Claude 3.7 的思考链分片 将推理思维过程作为一等公民,支持前端独立折叠框展示,与最终正文完全解耦。
ToolCallDeltaEvent index: int
delta: str
partial: Message
工具调用参数分片流式拼接过程 提供更细粒度的工具参数流式可观测性。
ToolCallDoneEvent index: int
tool_call: ToolCall
partial: Message
单个工具调用的 JSON 参数解析并验证完毕 提前通知调度层工具已完整就绪,使 UI 能在流式未完全结束前提前渲染“准备运行”态。
StreamDoneEvent message: Message
usage: dict[str, int]
底层网络流正常完型 交付合法定型的完整 Message,并携带真实 Token usage,为上下文压缩提供权威校准锚点。
StreamErrorEvent error: Message
stop_reason: str
exc: 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

# 1. 经典同步/异步接口(全部保持高度对称的干净签名)
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]: ...

# 2. Tau 对齐高阶事件流接口
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:适配器只需发射简单的 ProviderTextDeltaEventProviderThinkingDeltaEventProviderToolCallEvent
    • 外层事件(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 # 增量文本:"world"
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_events20~30 行,只有纯粹的 if isinstance 事件转译,一眼见底。
状态一致性 微内核在累加、UI 在累加、中间件在暂存,极易发生状态不一致或时序颠倒。 单一权威来源,模型层是唯一的装配车间,随路广播 partial 快照。
思维链支持 只能用各种 metadata 临时补丁到处拼接,展示与存储逻辑别扭。 原生支持 ThinkingDeltaEvent 与有序块,思维链是一等公民。
新模型接入 每个 Provider 都要考虑如何累加、如何处理异常与中断。 适配器只需抛出最简内部 delta,中央累加器自动赋能全量事件。
系统稳定性 网络波动或取消信号直接引发未捕获异常击穿 ReAct 循环。 Never-Throw 保证,失败与取消下沉为带规范 stop_reason 的实体,驱动后续自愈。

核心设计原则与架构不变式总结

  1. Never-Throw Guarantee(模型层崩溃免疫): 网络断开、超时或大模型产出畸形格式时,模型层绝不抛出未捕获异常打崩主程序,统一就地封装为带 stop_reason="error"StreamErrorEvent,引导上层自愈或安全闭环;
  2. Strict Timing Order Invariant(严密生命周期时序): 流式事件严格遵循 StreamStartEvent ➔ *DeltaEvent... ➔ StreamDoneEvent / StreamErrorEvent 的不可逆确定性时序,彻底消灭首字延迟前提前发射或结束消息倒置的时序 Bug;
  3. Structured First-Class Citizens(结构化实体一等公民): 参数在模型层完成反序列化,彻底消除 4 重 JSON 编解码死循环(4x JSON Ping-Pong),直连微内核调度与工具执行。