功能设计

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),直连微内核调度与工具执行。

生命周期事件

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Agent 生命周期
agent_start
agent_end(messages)

// Turn 生命周期(一轮 = 一次助手响应 + 工具调用/结果)
turn_start
turn_end(message, toolResults)

// 消息生命周期(user / assistant / toolResult 消息都会发)
message_start(message)
message_update(message) // 仅异步流式
message_end(message)

// 工具执行生命周期(ToolExecutionStart/End 可被 hook 干预)
tool_execution_start(toolCallId, toolName, args)
tool_execution_update(toolCallId, toolName, args, partialResult) // 仅异步流式
tool_execution_end(toolCallId, toolName, result, isError)

事件集对齐 pi 的生命周期模型(Agent/Turn/Message/Tool 四组,正常执行 start/end 成对;被拦截/畸形参数的调用不发射 End)。 每个事件自动带 timestamp(Unix 秒,实例化时刻)。

事件的作用

一、事件的作用:把「可观察」变成「可扩展」

一句话:事件机制让 Agent.run() 成为一个「可被插桩的黑盒」——循环本身写死了,但循环的每一步都广播出去,让任何新功能都能「监听而不改框架」。

先看现在 run() 里实际发生了什么:

1
2
3
4
5
6
7
8
9
10
11
12
13
Agent.run("东京天气")

├─ AgentStart
├─ MessageStart(user) → MessageEnd(user) # 你的问题进 transcript
├─ TurnStart(1)
│ ├─ MessageStart(assistant) → MessageEnd(assistant) # 模型想调工具
│ ├─ ToolExecutionStart → ToolExecutionEnd # 执行 get_weather
│ ├─ MessageStart(tool) → MessageEnd(tool) # 结果写回
├─ TurnEnd(message=assistant, tool_results=[...])
├─ TurnStart(2)
│ ├─ MessageStart(assistant) → MessageEnd(assistant) # 模型直接回答
│ └─ TurnEnd(message=assistant, tool_results=[]) # 本轮工作闭环结算
├─ AgentEnd(messages=[...], ...)

这 8 个事件(同步阶段)就是「循环的完整心电图」——每一段代码执行都对应一个事件。而事件的作用 = 把「这段代码发生了」告诉任何想知道的人。

二、核心:为什么事件 = 扩展点

关键洞察在于——你的所有后续功能,几乎都是「想观察循环在干什么」,而不是「想改循环本身」。事件机制让这些功能都变成「加一个 hook」,而不是「改 run()」。

看三个后续阶段怎么用:

扩展 1:session 落盘(阶段 3)—— Agent 显式双写(不是 hook 消费事件流)

1
2
3
4
# 实际实现:Agent.run() 里每加一条消息,同时写内存 + 落盘
# (不走 hook——session 是核心持久化,Agent 显式调用)
self.messages.append(user_msg) # 内存 transcript
self.session.add_message("user", user_input) # 树 + 原子落盘(session 必填,无条件)

(注:早期设计设想过”用 hook 消费事件流写 session”,但实际实现是 Agent 双写——session 是核心功能,显式调用比 hook 更直接、不依赖事件顺序。)

扩展 2:context 压缩(阶段 4)—— prepare 每轮调用(不是监听事件)

1
2
3
4
5
6
# 实际实现:prepare() 在每次 llm.chat 前调用,内部判断是否超阈
# (不走"监听 MessageEnd 计数"——压缩是发送前视图变换,时机在发消息前)
while ...:
view = self._ctx.prepare(self.messages) # 四层管线(超阈 → 摘要)
resp = self.llm.chat(messages=view, tools=...)
...

ContextCompacted 事件早在设计时就预留了(阶段 4 用)——这就是「事件先定义、后发射」的意义:接口提前定好,实现到时接上。

扩展 3:coding agent UI(未来)—— 渲染过程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# 新功能:TUI 渲染 —— 监听事件更新界面
def ui_render(event):
if isinstance(event, TurnStart):
ui.show_round(event.iteration)
elif isinstance(event, ToolExecutionStart):
ui.show_tool_call(event.tool_name, event.args)
elif isinstance(event, MessageEnd):
ui.show_message(event.message.content)
return None

```python
# 观察者只读订阅(UI 渲染、日志打点)
unsubscribe = agent.subscribe(ui_render)

# 决策拦截门禁(权限检查、参数改写)
agent = Agent(llm=llm, tools=tools, hooks=[
(UserInputHook, input_sanitizer),
(ToolCallHook, permission_guard),
(ToolResultHook, output_redactor),
])
```

hook系统

架构解耦重塑(对标 Pi & Tau):在最新架构中,只读事件(Event决策门禁(*Hook实现彻底的物理与概念解耦: - 只读事件(events.pyAgentStart, TurnEnd, MessageEnd, ToolExecutionStart 等事实通知,统一通过 agent.subscribe(listener) 进行只读旁路监听; - 决策门禁(hooks.pyUserInputHook, AgentStartHook, BeforeModelCallHook, ToolCallHook, ToolResultHook 5 大专职拦截载荷,在 Agent(hooks=[(HookClass, callback), ...]) 中批量注册,负责阻断、参数改写与结果篡改。

多 hook 注册表

核心两个函数:注册(register)和触发(emit)。

注册:把 hook 回调挂到 Hook 目标名下

1
2
def register(self, target: type, callback: Callable):
self._hooks.setdefault(target, []).append(callback)
  • target = Hook 类型(如 ToolCallHook
  • callback = 门禁函数((*Hook) -> HookResult | None
  • 同一个 Hook 类型可挂载多个回调,按注册顺序保序存储

触发:到点时异步逐个调用(支持 async/sync 混合回调)

1
2
3
4
5
6
7
8
9
10
11
12
async def emit(self, event: Any) -> HookResult | None:
"""异步触发 Hook,支持协程与普通函数,返回第一个非 None 结果(短路)。"""
for cb in self._hooks.get(type(event), []):
try:
res = cb(event)
if inspect.isawaitable(res):
res = await res
if res is not None:
return res # 短路:第一个干预结果即刻生效返回
except Exception as exc:
logger.error(f"Hook callback {cb} error: {exc}", exc_info=True)
return None

HookResult 干预模型与五大决策拦截点

Hook 回调返回 None = 放行;返回 HookResult = 强力干预。

1
2
3
4
5
6
7
8
9
10
@dataclass(frozen=True)
class HookResult:
block: bool = False # 是否阻断当前执行
reason: str | None = None # block 阻断时的原因说明
terminate: bool | None = None # 三态:是否请求微内核提前退出当前批次
updated_input: str | None = None # 改写用户原始输入 (UserInputHook 用)
updated_system_prompt: str | None = None # 改写系统提示词 (AgentStartHook 用)
updated_messages: list | None = None # 临时改写送入模型的上下文 (BeforeModelCallHook 用)
updated_args: dict | None = None # 改写工具入参 (ToolCallHook 用)
updated_result: str | None = None # 改写工具执行结果 (ToolResultHook 用)

五大独立生命周期决策拦截点(hooks.py

  1. UserInputHook:截获用户原始输入,支持 block 阻断或 updated_input 前置清洗与改写;
  2. AgentStartHook:启动前拦截,支持 updated_system_prompt 动态更新首条 system 消息;
  3. BeforeModelCallHook:调 LLM 前拦截,支持 updated_messages 临时改写视图(临时 View 改写 vs 真实 Session 零污染);
  4. ToolCallHook:工具执行前拦截,支持 block 拦截高危命令或 updated_args 动态修补参数;
  5. ToolResultHook:工具执行后拦截,支持 updated_result 篡改出参,以及 terminate 提前退出整批 ReAct 循环。

用起来的样子

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
# 注册阶段(应用层,构造 Agent 后)
def guard_hook(event) -> HookResult | None:
"""权限门控(挂到 ToolExecutionStart,可拦截/改参数)"""
if isinstance(event, ToolExecutionStart):
if event.tool_name == "bash" and "rm -rf" in event.args.get("command", ""):
return HookResult(block=True, reason="危险命令被拦截")
if event.tool_name == "write_file":
return HookResult(updated_args={**event.args, "path": abs(event.args["path"])})
return None

def log_hook(event) -> HookResult | None:
"""日志(挂到 ToolExecutionStart,纯观察)"""
if isinstance(event, ToolExecutionStart):
print(f"[hook] {event.tool_name}({event.args})")
return None

agent = Agent(llm=llm, tools=tools, hooks=[
(ToolExecutionStart, guard_hook), # 同一事件挂多个
(ToolExecutionStart, log_hook),
])

机制一句话:注册表 = 事件类 → 回调列表;触发 = 到点时按序调用,短路返回。 你说「在某个事件时调用注册的 hook」——完全正确,_emit 就是那个「到点」的调用点。

Agent类(轻量 Harness 架构演进)

在经历 Tau 对齐重塑(阶段 18)后,Agent 类已经从早期“状态+循环+调度揉在一个类”的上帝类,全面演进为对标 Tau / Pi 的轻量有状态外壳(AgentHarness

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
┌─────────────────────────────────────────────────────────────────┐
│ class Agent (AgentHarness) │
│ "轻量宿主外壳:持有状态 + 暴露原生事件流" │
│ │
│ ┌───────────── __init__ 装配 ─────────────┐ │
│ │ llm ← 注入(模型边界) │ │
│ │ session ← 会话(只追加存储驱动) │ │
│ │ registry ← 工具注册表 │ │
│ │ messages ← 状态/记忆(自愈清洗) │ │
│ │ message_queue ← 干预队列 (Steer/Followup) │
│ │ hooks ← HookRegistry(拦截管线)│ │
│ └─────────────────────────────────────────┘ │
│ │
│ ┌────────── 一等公民核心驱动:prompt_stream() ──────────┐ │
│ │ 异步生成器:run_agent_loop(llm, tools, messages...) │ │
│ │ │ │ │
│ │ ▼ 逐一 yield 生命周期事件流 │ │
│ │ AgentStart ➔ TurnStart ➔ MessageUpdate ➔ TurnEnd... │ │
│ └───────────────────────────────────────────────────────┘ │
│ │
│ ┌────────── 经典便利门面:run() ──────────┐ │
│ │ async def run(user_input): │ │
│ │ async for event in prompt_stream: │ │
│ │ return final_text │ │
│ └────────────────────────────────────────┘ │
│ │
│ ┌────────── 观察者监听:subscribe() ──────┐ │
│ │ subscribe(listener) -> unsubscribe │ │
│ └────────────────────────────────────────┘ │
│ │
│ ┌────────── 外部纯函数微内核 (loop.py) ───┐ │
│ │ run_agent_loop: 彻底接管两阶段调度与双层循环 │
│ │ _provider_context: 前置剥离空异常 + 自愈断头 │
│ └─────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘

对应代码:

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
class Agent:
def __init__(
self,
*,
llm: LLM,
tools: list[Tool],
session: Session,
system_prompt: str | None = None,
max_iterations: int | None = None,
context_budget: int | None = None,
hooks: list[tuple[type[Event], Callable]] | None = None,
...
): ...

# 1. 现代化原生事件流生成器(面向 TUI / CLI 打字机 / 前端)
async def prompt_stream(self, user_input: str) -> AsyncIterator[Event]:
"""将用户输入送入微内核,实时产生事件流,同时驱动 session 落盘。"""
# 前置决策点 UserInput / AgentStart ...
async for event in run_agent_loop(
llm=self.llm,
model=self.model,
system=system_prompt,
messages=self.messages,
tools=self.registry,
context_manager=self._ctx,
get_steering_messages=self._get_steering_messages,
get_follow_up_messages=self._get_follow_up_messages,
before_model_call=self.hooks.emit,
before_tool_call=self.hooks.emit,
after_tool_call=self.hooks.emit,
):
# 同步落盘与外部观察者广播
yield event

# 2. 经典便利门面(100% 向后兼容老脚本与测试)
async def run(self, user_input: str) -> str | None:
"""运行并返回最终文本(内部消费 prompt_stream)。"""
final_text = None
async for event in self.prompt_stream(user_input):
if isinstance(event, AgentEnd):
final_text = event.final_text
return final_text

# 3. 观察者订阅 API
def subscribe(self, listener: Callable[[Event], None]) -> Callable[[], None]:
"""订阅事件流(返回注销句柄)。"""
self._subscribers.append(listener)
return lambda: self._subscribers.remove(listener)

def abort(self) -> None: ... # 协作式取消当前任务并自愈落盘
def steer(self, message: str) -> None: ... # 即时转向
def follow_up(self, message: str) -> None: ... # 排队追加

三、调度总装枢纽:prompt_stream 全流程剖析与 run() 职责解构

在现代 Agent 架构中,对外暴露的驱动接口形成了清晰的分工: - async def prompt_stream(...):系统最高阶的一等公民事件流生成器(Producer & Orchestrator); - async def run(...):仅 10 行代码的阻塞式便利消费门面(Pure Consumer)

1. prompt_stream 的六阶段全生命周期流水线

当外部调用 async for event in agent.prompt_stream("帮我重构代码"): 时,它严格按以下 6 大阶段推进:

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
用户输入 user_input


【阶段 1: 状态复位与扩展懒加载】
└─ 复位 _aborted = False,生成新的 CancellationToken 取消令牌
└─ 若插件/扩展尚未加载,触发 self.extension_manager.load()


【阶段 2: 用户输入安全门禁 (Hook 1: UserInputHook)】
└─ 审查输入是否违规或需脱敏:
├─ block=True ──> 发射 AgentEnd(stop_reason="blocked") 并直接 return 退出!
└─ 改写输入 ────> 更新 user_input 为脱敏后的新文本


【阶段 3: 会话树同步与转录本拓扑自愈】
└─ 读取当前 Session 路径历史:restored = system + get_full_history_messages()
└─ 运行 repair_tool_history(restored) 瞬间缝合所有断头调用,免疫 API 400


【阶段 4: 启动提示词门禁 (Hook 2: AgentStartHook)】
└─ 审查/改写 System Prompt (允许插件动态植入全局规则)


【阶段 5: 微内核点火启动 (run_agent_loop)】
└─ 将准备好的参数喂给 run_agent_loop,获得异步生成器 loop_gen


【阶段 6: 三向分发大循环 (async for event in loop_gen)】
├─ 1. 磁盘持久化分支 ──> 监听到 MessageEnd ➔ self.session.add_message 落盘
├─ 2. 旁路广播分支 ──> 遍历 self._subscribers 广播给第三方监听器
└─ 3. 主干流式分支 ──> yield event 实时推给外部调用方

逐阶段技术细节解密:

  1. 阶段 1(状态复位与扩展懒加载)
    为本次调用创建专属的 CancellationToken。采用懒加载策略:在第一次真正发起问答时才执行 await self.extension_manager.load(),保证 Agent 初始化零异步卡顿;
  2. 阶段 2(用户输入门禁 UserInputHook
    防守前置,零成本拦截。在用户输入触碰模型或写入历史前先行安检。若插件判定违规(block=True),直接就地发射 AgentEnd(stop_reason="blocked") 退出,不花一分钱模型费用;若需要脱敏(updated_input),则无缝篡改后放行;
  3. 阶段 3(会话树指针同步与转录本自愈)
    • 解决 /rewind(会话回滚)后的指针同步问题,严格以磁盘 Session 树当前指针为基准重新拉齐内存历史;
    • 运行 repair_tool_history,即使上次任务中途被强制杀死或崩溃,残留的断头工具调用也能被毫秒级修复,绝对杜绝大模型 API 报 400 Bad Request
  4. 阶段 4(系统提示词门禁 AgentStartHook
    允许策略插件根据当前用户身份或环境,动态向 system_prompt 中追加安全规约;
  5. 阶段 5(微内核点火启动)
    将组装完毕的消息历史、工具注册表、取消信号和队列回调打包,调用 run_agent_loop 启动纯函数引擎;
  6. 阶段 6(三向分发大循环 Tee Splitter)
    微内核吐出的每一个事件,在此处被“一鱼三吃”
    • 向磁盘写:一旦捕获到 MessageEnd(说明消息已完整终态定型),立刻调用 self.session.add_message 通过原子重命名安全落盘到 .jsonl
    • 向广播发:遍历 self._subscribers 列表,安全通知所有旁路监听器(日志、监控、状态栏);
    • 向调用方推:通过 yield event 原生推向终端屏幕或 WebSocket。

2. prompt_stream 与 run() 的职责解构

很多传统框架的代码之所以臃肿,是因为把“流式输出”和“最终结果返回”混在了同一个方法里。而在现代微内核架构中,两者被进行了极其清晰的“生产者 vs 消费者”职责解构:

维度 prompt_stream(...) run(...)
返回类型 AsyncIterator[Event](异步生成器事件流) str | None(最终纯文本答案)
角色定位 调度总装者与事件生产者(Orchestrator & Producer) 端到端纯消费者(Pure Consumer)
内部职责 负责安检门禁、会话同步、驱动微内核、写磁盘、旁路广播 完全不关心底层调度,只管迭代 prompt_stream 等到终态
面向场景 终端 TUI 彩色打字机、WebSocket 实时推送、长程任务进度监听 离线单测(0.001 秒断言输出)、自动化脚本批处理、极简命令行调用
异常体现 将错误包装为结构化 StreamErrorEventAgentEnd(stop_reason="error") 发射 遇到 stop_reason == "error" 时主动向上抛出 RuntimeError,遇到取消返回 "(cancelled)"

极简消费门面:run() 的优雅退化

正是因为 prompt_stream 把所有的脏活累活(安全拦截、自愈、驱动微内核、写文件)全部在底层包揽了,上层的便利方法 run() 才能做到前所未有的纯粹与轻薄——它退化成了仅仅 10 行代码的普通事件消费者

1
2
3
4
5
6
7
8
9
10
11
12
13
async def run(self, user_input: str) -> str | None:
"""追加 user 消息 → 内部消费 prompt_stream 事件流 → 返回最终文本。"""
final_text = None
async for event in self.prompt_stream(user_input):
if isinstance(event, AgentEnd):
if event.stop_reason == "cancelled":
return "(cancelled)"
if event.stop_reason == "blocked":
return event.final_text or "(blocked)"
if event.stop_reason == "error":
raise RuntimeError(event.final_text or "Error during model stream")
final_text = event.final_text
return final_text

这种职责解构的收益极为巨大: 1. 单一事实来源(Single Source of Truth):整个框架只有 prompt_stream 一条真正的执行流水线,任何安全补丁、持久化优化或时序更新只需改动一处,run() 自动享受到所有红利; 2. 兼得双重体验:既能满足现代交互界面对毫秒级流式响应的严苛要求,又能保留单测和简单脚本中 ans = await agent.run(...) 一行调用的极致清爽!

3. 只读事件订阅管道:subscribe 与 _notify 广播器

在生产级 Agent 中,UI 状态更新、日志埋点、监控探针等外部系统需要实时感知 Agent 状态,但绝不能侵入或打断核心推理流程Agent 类实现了对标 Pi(agent-session.js)与 Tau(harness.py)的只读轻量观察者管道:

1
2
3
4
5
# 订阅生命周期事件流(返回注销句柄)
unsubscribe = agent.subscribe(my_ui_listener)

# 业务注销(消除内存泄漏隐患)
unsubscribe()

_notify 广播器的六大工程保证:

1
2
3
4
5
6
7
async def _notify(self, event: Event) -> None:
"""将生命周期事件安全广播给所有旁路订阅者(对标 Tau AgentHarness._notify)。"""
for sub in list(self._subscribers):
with contextlib.suppress(Exception):
res = sub(event)
if inspect.isawaitable(res):
await res
  1. 快照遍历(Snapshot Iteration)list(self._subscribers) 保证在迭代过程中即使外部调用了 unsubscribe() 也绝不会引发并发修改异常;
  2. 异常完全隔离(Never-Throw Guarantee):外部监听器抛出的任何异常被 contextlib.suppress(Exception) 彻底吞默,第三方的 bug 绝对无法击垮 Agent 主循环;
  3. 同异步无缝自适应:自动检查 inspect.isawaitable(res),无论监听器写成普通同步函数还是 async def 异步协程,框架均能自适应等待;
  4. 闭包注销句柄(Teardown Token):解耦索引,注销时按对象引用精准剔除;
  5. 全量关键触点覆盖:在用户输入拦截阻断、系统提示词阻断、微内核主事件迭代循环、以及上下文压缩回写等所有生命周期转折点统一触发广播。

核心架构演进:为什么要把 run_agent_loop 抽离为独立纯函数?

在阶段 18 对标 Tau 架构的重塑中,我们做出了一个决定性的重构动作:将原先内联在 Agent 类中长达 260 行的双层循环彻底剥离,提炼为 packages/my-agent-core/src/my_agent_core/loop.py 中的独立纯函数生成器 run_agent_loop

这一解耦设计带来了两大根本性的架构质变:

1. 纯函数无状态微内核的威力(五大工程质变)

传统框架喜欢把循环写死在类方法内部(class Agent: def run(self): while True:),导致循环与类状态深度咬合。在长对话或复杂场景下,这种设计弊端尽显。而抽离为独立的 async def run_agent_loop(...) 后,系统获得了五大不可替代的红利:

  1. 彻底消除隐式状态突变(零 self,绝对确定性)
    • 在类方法中,循环随时可能修改 self.xxx,多轮对话后状态满天飞;
    • 纯函数微内核完全没有 self,零隐式状态!所有输入项(大模型门面 llm、消息列表 messages、工具注册表 tools、取消令牌 signal)全部通过参数显式传入。输入确定则输出事件流绝对确定,排查调试的确定性拉满。
  2. 通用计算发动机:多场景极致复用(Write Once, Run Everywhere)
    • run_agent_loop 就像一台通用的汽车发动机,不绑定任何特定底盘,能够随处无缝挂载:
      • 交互式 CLI:配上 Agent 外壳,驱动 Session 树持久化与终端彩色打字机;
      • Web / API 流式服务:直接挂载到 FastAPI 的 StreamingResponse 或 WebSocket,中间事件流无需中间人直接推送到前端;
      • 大规模无头评测(Headless Benchmarks):跑 SWE-bench 评测 1000 道题时,无需任何 Session 落盘与复杂插件装配,直接拿微内核裸跑,内存开销降低 80%,极速并发;
      • Subagent 子代理:轻量级临时执行手脚任务,无需背负庞大的 Agent 实例。
  3. 单元测试速度提升 100 倍(零装配、纯内存秒跑)
    • 以前测试循环逻辑,必须先装配一整台 Agent 大机器(Session 树、MemoryStore、Skills、Plugins),单测准备代码冗长且慢;
    • 现在在 tests/test_agent_loop_pure.py 中,直接传入 FakeLLM 和列表,0.001 秒测完并发工具与动态转向(11 项复杂微内核单测在 0.05 秒内全通)。
  4. 动静分离,职责正交(SRP 原则)
    • agent.py(类)只负责静态的结构与持久化(持有 Session 树、长期记忆、扩展装配);
    • loop.py(函数)只负责动态的过程推理与调度(怎么推理、怎么并发调工具、怎么转向、怎么中断)。
  5. 双层循环拓扑与动态干预毫秒级响应
    • 内层微观循环在工具执行间隙动态收割 get_steering_messages() 实现即时转向;
    • 外层宏观循环在轮次自然收尾处收割 get_follow_up_messages() 自动衔接追问任务。

2. 将中间“事件流”提升为系统级的一等公民(First-Class Stream)

大模型的运行本质上是一个基于时间轴的渐进流式过程:Token 在一个字一个字蹦出、思考块在逐步形成、工具在按需触发并耗时运行。

传统 Callback / Hook 的三大灾难(打洞插管):

在传统黑盒设计中,函数只能最后 return final_text。外部若想知道中间状态,必须在外面定义一堆回调函数(on_tokenon_tool_start)硬塞给 Agent。 - 控制权倒置:调用方不能主动拉取数据,只能被框架被动调用; - 异常黑洞:回调函数一旦抛出异常,会直接击垮主循环,或者被框架静默吞掉; - 无法优雅暂停与取消:外部用户在前端点了停止,回调依然在疯狂发射,根本停不下来。

生成器(AsyncIterator[Event])的一等公民之美:

在现代 run_agent_loop 架构下,事件流就是函数的正常返回值本身!没有任何侧路管子,没有控制权倒置!

  1. 调用方掌握 100% 绝对主动权(传送带模型): 外部调用者只需写出自然优雅的 async for 循环,站在传送带旁边接住事件:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    # 外部使用代码极度自然、主动:
    async for event in agent.prompt_stream("帮我计算 37 * 19"):
    if isinstance(event, MessageUpdate):
    print(event.chunk.delta, end="", flush=True) # 原生打字机!
    elif isinstance(event, ToolExecutionStart):
    print(f"\n[正在调用工具: {event.tool_name}]")
    elif isinstance(event, ToolExecutionEnd):
    print(f"\n[工具执行完毕: {event.result}]")

    # 想停就停:直接 break 即可触发生成器安全清理与底层中断自愈!
    if user_clicked_stop:
    break
  2. 让上层门面实现前所未有的轻薄: 因为事件流成了一等公民,外层的 agent.run() 彻底被掏空,退化成了仅 10 行代码的纯事件消费者:
    1
    2
    3
    4
    5
    6
    7
    async def run(self, user_input: str) -> str | None:
    """自己只做普通消费者,消费事件流并返回最终文本。"""
    final_text = None
    async for event in self.prompt_stream(user_input):
    if isinstance(event, AgentEnd):
    final_text = event.final_text
    return final_text
    一套微内核底座,同时完美满足了“前端极致细粒度的响应式流渲染”与“后端一行代码调用的极简体验”!

3. 架构深思:Pi 的 ExtensionAPI 如何将“只读事件”与“拦截决策点”相辅相成?

在深入对比 Pi(@earendil-works/pi-coding-agent)与 Tau 的设计后,我们终于看清了困扰很多开发者的核心谜团:为什么在很多初学者的 Agent 代码里,事件和 Hook 会被写成一团乱麻?而 Pi 却能写得如此规整有序?

答案就在于:Pi 深刻认清了“只读事件(传信息)”与“决策拦截点(改控制流)”是两套完全正交的体系,并通过 ExtensionAPI 的三层架构实现了二者的完美相辅相成!

① 概念正交化:纯只读事件 vs 决策拦截点

体系类别 代表事件名称 核心职责 返回值意义 典型消费场景
纯只读事件流
(Agent Events)
turn_start, turn_end, message_update, tool_execution_end 客观传递信息(Stream of Facts),单向广播发生了什么 被动忽略 (void),没有资格改变控制流 驱动终端 TUI 彩色打字机、状态栏、Web 前端进度条、日志审计
五大决策拦截点
(Decision Points)
input, before_agent_start, context, tool_call, tool_result 主观裁决与改写(Control Gates),把守关键命脉卡口 决定生死!返回 { block: true } 阻断或篡改参数/上下文 拦截 rm -rf 高危命令、人机审批流、动态注入提示词、清洗脱敏敏感词

② 为什么初学者(包括我们重构初期)会陷入杂乱?

  • 概念混淆的陷阱:我们把所有的只读事件和拦截点硬塞进了同一个 HookRegistry,导致微内核里充斥着几十处毫无意义的 if hook_registry: await hook_registry.emit(ev)
  • 双通道视觉噪音:每个事件既要 yield 给外部,又要 await emit 给内部。这就好比广播电台播报“天亮了”(纯通知),却非要在广播室门口设个安检站一本正经地问安检员“天亮了需要拦截吗?”。

③ Pi 的三层解耦美学:ExtensionAPI 与微内核的精密咬合

Pi 极为高明地通过三层架构化解了这一矛盾:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
┌─────────────────────────────────────────────────────────────┐
│ 【外部插件层 (Extension Developer)】 │
│ 面向开发者提供极简统一入口:pi.on(event, handler) │
│ • 监听只读事件:pi.on("message_update", handler) │
│ • 拦截安全门禁:pi.on("tool_call", handler) │
└──────────────────────────────┬──────────────────────────────┘
│ 智能分流

┌─────────────────────────────────────────────────────────────┐
│ 【中间调度层 (ExtensionManager)】 │
│ • 只读事件 ➔ 接入单向广播观察者队列 (Observer Queue) │
│ • 决策拦截点 ➔ 串联为洋葱模型中间件流水线 (Middleware) │
└──────────────────────────────┬──────────────────────────────┘
│ 注入纯函数插槽

┌─────────────────────────────────────────────────────────────┐
│ 【底层微内核 (runAgentLoop)】 │
│ • 微内核对 ExtensionAPI 零感知!不持有全局 Hook 注册表 │
│ • 纯只读事件:微内核只管一行代码 yield event (纯单向推流) │
│ • 五大决策点:微内核仅通过入参暴露 5 个具体干净的回调插槽 │
│ (beforeToolCall, afterToolCall, transformContext...) │
└─────────────────────────────────────────────────────────────┘

这种相辅相成的设计妙在哪里?

  1. 对插件作者(API 面)极致的直觉与统一。开发者不需要记两套完全不同的 API,只要在同一个 pi.on 后面写函数,TypeScript 类型系统会自动约束当前事件“能不能拦截、可以返回什么”;
  2. 对中间层(调度面)职责清晰的智能分流。只读事件走广播,安全门禁走洋葱流水线(支持多插件按序改写参数与一票短路阻断);
  3. 对微内核(底层引擎)极致的轻盈与规整。微内核根本不需要感知插件的存在,也不需要充斥 30 处 Hook 检查样板代码;它只在 5 个固定的命脉节点上留出回调插槽,其余时间顺风顺水地跑自己的纯状态机!

这正是顶级工业级 Agent 框架的架构魅力所在:上层 API 极其人性统一,底层微内核极其纯粹高效,中间通过清晰的分流管道相辅相成!

4. 旁路广播总线:subscribe() 订阅机制的底层本质与实战全景

很多开发者在接触到 agent.prompt_stream(...) 后会产生一个经典困惑: > “既然我已经能通过 async for event in agent.prompt_stream(...) 直接拿到所有事件了,为什么还需要 agent.subscribe(...) 订阅机制?这不是多此一举吗?”

答案就在于:prompt_stream 是单通道的“独食吸管”,而 subscribe 是多方共享的“公共大喇叭”。

① 为什么必须要有 subscribe?(两大刚需场景)

  1. 配合最常用的极简门面 await agent.run(): 大部分日常业务或测试中,我们只想要最终答案:ans = await agent.run("计算 37*19")。 然而,run() 内部已经把 prompt_stream消费吞掉了,外部只能拿到最终字符串。如果你此时既想用清爽的 agent.run(),又想在控制台打印工具调用日志,或者在状态栏转圈,唯一的办法就是提前挂载 agent.subscribe(my_listener)
  2. 克服 Python 生成器的“单消费者限制”: Python 的异步生成器不能被多个 async for 同时读取。如果你的系统需要同时做三件事:
    • 主流程:在终端打印大模型的话;
    • 日志系统:把所有工具调用记录进 agent.log
    • 监控系统:把消耗的 Token 上报云端; 若没有 subscribe,你必须把日志和监控代码全揉进主流程的 async for 里,代码变成大泥球。有了 subscribe,日志和监控模块只需各自独立 agent.subscribe(...)旁路静默监听,与主业务彻底解耦

② agent.subscribe(…) 的入参到底是什么?

入参是一个“函数”(由你编写的回调函数 Callable)。
该函数必须接收一个参数(框架会把发生的 event 实体传给它)。支持三种形式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 形式 1: 普通函数(最直观常用)
def my_listener(event: Event):
print("收到事件:", event)

agent.subscribe(my_listener) # 传入函数名,切勿带括号!

# 形式 2: 异步函数(适合写数据库或发网络请求)
async def my_async_listener(event: Event):
await websocket.send_json({"event": type(event).__name__})

agent.subscribe(my_async_listener)

# 形式 3: 匿名函数(适合一行调试)
agent.subscribe(lambda ev: logger.info("Event: %s", ev))

③ 订阅只能看到事件类型吗?(绝对不是!)

每个推送到订阅函数的 event,都是一个满载真实业务数据的强类型实体对象,你能从中拿到所有的具体数据:

  • 工具启动(ToolExecutionStart
    • event.tool_name:工具名(如 "bash""read_file");
    • event.args大模型传入的真实入参字典(如 {"path": "src/app.py"});
    • event.tool_call_id:工具调用唯一 ID。
  • 工具结束(ToolExecutionEnd
    • event.result工具执行的完整输出内容(如文件全文、命令执行结果);
    • event.is_error:执行是否出错(True / False);
    • event.terminate:是否触发了阶段 7 提前熔断标记。
  • 大模型流式吐字(MessageUpdate
    • event.delta刚刚吐出的那一个字/词(驱动终端流式打字机)。
  • 任务终结(AgentEnd
    • event.final_text:最终交付给用户的完整回答文本;
    • event.iterations:总共跑了多少轮推理;
    • event.stop_reason:为何结束("end_turn""max_iterations" 等)。
  • 上下文压缩(ContextCompacted
    • event.tokens_beforeevent.tokens_after:压缩前后的真实 Token 水位。

实战中,你可以轻松写出一个赏心悦目的终端监听器:

1
2
3
4
5
6
7
8
9
10
11
def pretty_printer(event: Event):
if isinstance(event, ToolExecutionStart):
print(f"🔧 [工具启动] {event.tool_name},参数: {event.args}")
elif isinstance(event, ToolExecutionEnd):
status = "❌ 失败" if event.is_error else "✅ 成功"
print(f"{status} [工具完成] 结果: {event.result}")
elif isinstance(event, MessageUpdate):
print(event.delta, end="", flush=True)

agent.subscribe(pretty_printer)
await agent.run("计算 37 * 19")

④ 底层实现原理(三部曲机制)

subscribe 在底层代码中仅有约 15 行,设计极其紧凑:

1
2
3
4
5
6
7
8
1. 登记通讯录 (subscribe)
agent.subscribe(listener) ──> self._subscribers.append(listener)

2. 逐个打电话广播 (prompt_stream 内部)
微内核产生新事件 ──> for sub in list(self._subscribers): sub(event)

3. 随时退订除名 (unsubscribe 闭包)
调用返回的 unsub() ──> self._subscribers.remove(listener)
  1. 闭包式退订令牌(Teardown Token)
    subscribe 返回一个私有闭包 unsubscribe()。调用方无需持有监听器列表或索引,直接调用 unsub() 即可瞬间安全退订,从机制上彻底杜绝内存泄漏
  2. list(self._subscribers) 浅拷贝快照
    在广播循环中使用快照遍历,防止监听器在执行过程中自行动态注销导致 RuntimeError: list changed size during iteration 崩溃;
  3. contextlib.suppress(Exception) 异常隔离
    第三方订阅函数内部若抛出任何未捕获异常,全部被静默隔离,绝对不影响 Agent 主调度循环的稳定性
  4. 同步/异步自适应(isawaitable
    底层通过 inspect.isawaitable(res) 自动检测返回值,若订阅函数是 async def,主线程自动进行 await 调度。

5. 微内核调度流水线独立专文导航

至此,Agent 类已经蜕变成一个高度轻量、专注会话状态管理与全系统依赖注入的 Harness 宿主外壳。

关于微内核内部的双层状态机循环(宏观 Follow-up 追问 + 微观 ReAct 转向)、单轮 8 步时序流水线、大模型推理车间(_assistant_turn)以及工业级工具执行车间(_execute_tools_turn)的深度解析,已独立拆分并收录于专属博客中:

👉 《my-pi-agent–loop微内核.md》

前言

为了更深入理解agent的工程实现,本文会逐步从底层搭建一个agent(不借助任何框架如langchain等),agent是在循环中调用工具的模型,直到给定任务完成(ReAct架构),如下图所示。

image-20260731213620331

核心组件如下

image-20260731214119395

工具系统

先纠正一个常见误解

模型是怎么知道”该调工具了”还是”该直接输出答案”的?这个判断是 LangChain/LangGraph 实现的吗?

“判断”根本不是 LangChain/LangGraph 实现的,是模型本身的能力。

OpenAI 等厂商对模型做过 function calling(工具调用)专项训练:模型学会了一件事——当请求里带有工具描述、且对话内容需要工具时,输出一个结构化的 tool_calls;不需要时,输出普通文本。这个决策发生在 OpenAI 服务器上的模型推理过程中,LangChain 源码里没有、也不可能有一行”决定何时调工具”的逻辑。

工具调用流程

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
=== 1. 上行翻译:Python 函数 -> 模型看得懂的 JSON schema ===
{
"type": "function",
"function": {
"name": "multiply", ← 来自函数名
"description": "Multiply two integers.", ← 来自 docstring
"parameters": { ← 来自类型标注 a: int, b: int
"properties": {
"a": { "type": "integer" },
"b": { "type": "integer" }
},
"required": ["a", "b"],
"type": "object"
}
}
}

=== 2. 实际发给 OpenAI 的完整请求 payload ===
{
"model": "gpt-4.1-mini",
"stream": false,
"tools": [ ...上面那段 schema... ], ← tools 是请求的顶层字段
"messages": [
{ "content": "Use the multiply tool to calculate 37 times 19.", "role": "user" }
]
}

=== 3. 下行翻译:OpenAI 原始响应 -> AIMessage.tool_calls ===
message 类型: AIMessage
content: ''
tool_calls: [{'name': 'multiply', 'args': {'a': 37, 'b': 19}, 'id': 'call_abc123', 'type': 'tool_call'}]
路由判断 bool(msg.tool_calls) = True -> 去 tools 节点

=== 4. 模型决定直接回答(不调工具)时 ===
tool_calls: []
bool(msg2.tool_calls) = False -> 去 __end__
image-20260731221318443

工具参数

人类选择工具前需了解工具的功能、使用场景和输入参数。大模型同理——模型依据这些信息选择合适的工具。按以下JSON格式提供工具信息。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
{
"type": "function",
"function": {
"name": "get_current_weather",
"description": "当你想查询指定城市的天气时非常有用。",
"parameters": {
"type": "object",
"properties": {
"location": {
"type": "string",
"description": "城市或县区,比如北京市、杭州市、余杭区等。"
}
},
"required": ["location"]
}
}
}
  • type字段固定为"function"
  • function字段为 Object 类型;
    • name字段为自定义的工具函数名称,建议使用与函数相同的名称,如get_current_weatherget_current_time
    • description字段是对工具函数功能的描述,大模型会参考该字段来选择是否使用该工具函数。
    • parameters字段是对工具函数入参的描述,类型是 Object ,大模型会参考该字段来进行入参的提取。如果工具函数不需要输入参数,则无需指定parameters参数。
      • type字段固定为"object"
      • properties字段描述了入参的名称、数据类型与描述,为 Object 类型,Key 值为入参的名称,Value 值为入参的数据类型与描述;
      • required字段指定哪些参数为必填项,为 Array 类型。

发起 Function Calling 前,在代码中定义工具信息数组(tools),包含每个工具的函数名、描述和参数定义。该数组在后续请求时作为参数传入。

工具注册机制

把「一个函数变成可调用的工具」其实分两步:

  1. @tool(打包):把一个函数 → 一个 Tool(说明书 + 函数本体)。产出的是一个对象。
  2. ToolRegistry.register(注册):把一堆 Tool 收进注册表(内部 name → Tool 字典)。产出的才是一张能按名查人的花名册,registry.execute 就靠它,凭模型回传的名字字符串反查到真函数。

函数注册为工具

1
2
3
4
@tool
def multiply(a: int, b: int) -> int:
"""Multiply two integers."""
return a * b

这段和下面完全等价:

1
2
3
4
5
def multiply(a: int, b: int) -> int:
"""Multiply two integers."""
return a * b

multiply = tool(multiply) # ← @tool 就是帮你写了这一行

所以 tool 就是个普通函数:吃进去一个函数,吐出来一个 Tool。@ 只是语法糖。

这里藏着一个你最好亲自验证一下的事实——执行完 multiply = tool(multiply) 之后,multiply 这个名字已经不是函数了,而是一个 Tool 对象。原来的函数被塞进了 Tool.func 里存着。

1
2
3
4
5
6
class Tool:
"""一个可被模型调用的工具:函数本体 + 发给模型的 JSON schema(类化重构后)。"""
func: Callable[..., Any] # 函数本体
name: str # 工具名(默认函数名)
description: str # 描述(默认 docstring)
params_model: type[BaseModel] # pydantic 参数模型(create_model 动态建模)

Callable[..., Any] 是什么

Callable 来自 from typing import Callable(第 12 行),它是一个泛型类型,写法是 Callable[[参数类型...], 返回类型]。比如:

1
2
3
Callable[[int, str], bool]
# └──┬───┘ └┬┘
# 接收 int和str 返回 bool

完整调用链

第 1 步:装饰时把函数存进 Tool

1
2
3
4
5
6
return Tool(
name=func.__name__, # 比如 "get_weather"
...
func=func, # ← 原函数本体存进来了
)

第 2 步:调用时按名字找回这个包裹(registry.execute 内部)

1
2
name = tc["function"]["name"]           # 模型说:"我要调 get_weather"
target = registry.get(name) # 查注册表,拿到对应的 Tool 对象(没有则 None)

第 3 步:解析模型给的参数(registry.execute 内部)

1
2
3
args = json.loads(tool_call.function.arguments)
# 模型传来的是 JSON 字符串,比如 '{"city": "北京"}'
# 解析后变成 Python dict:{"city": "北京"}

第 4 步:用 func 真正调用(Tool.execute 内部)

1
2
result = target.execute(args)           # 校验 + 执行(pydantic 参数校验,永不抛)
# 等价于:get_weather(city="北京")

**args 是字典解包,把 {“city”: “北京”} 展开成关键字参数 city=“北京” 传给 func。这就是「利用 Tool 对象的 func 调用函数」的确切时刻。

create_model 动态建模工具参数

为什么必须动态?

框架是「库」,不知道你会写什么工具

my_agent_core 是被 main.py 使用的库。库的代码在写的时候,根本不知道使用者会注册哪些工具:

1
2
3
4
5
6
7
8
@tool
def get_weather(city: str) -> str: ... # 用户可能写 1 个字段

@tool
def multiply(a: int, b: int) -> int: ... # 可能写 2 个字段

@tool
def search_docs(query: str, tags: list[str], limit: int = 5) -> str: ... # 3 个字段,类型各异

每个工具的参数形状都不一样。如果模型类是静态的,框架就得在源码里把「所有可能的工具签名」都写成类——那是不可能的。唯一的出路是:模型类在运行时、根据实际收到的函数来造。

实现

1
2
3
4
5
6
# tools.py:76
model = create_model(
f"{func.__name__}_Args", # 类名,如 "get_weather_Args"
__config__=ConfigDict(extra="forbid"),
**fields, # ← 关键
)

**fields 把 {“city”: (str, …)} 展开成关键字参数,等价于直接写:

create_model("get_weather_Args", __config__=..., city=(str, ...))

schema 在完整闭环里的角色

1
2
3
4
5
6
7
8
① 框架 → 模型:发送 schema("我有这些工具,参数格式如下")

② 模型 → 框架:返回 tool_call
name: "multiply"
arguments: '{"a": 37, "b": 19}' ← 模型按契约生成的合规参数

③ 框架执行:json.loads(arguments) → {"a": 37, "b": 19}
func(**args) → 703 ← registry.execute 内部

第 ② 步值得多看一眼:‘{“a”: 37, “b”: 19}’ 这个 JSON 字符串是模型自己生成的——它读了 schema,知道该给 a 和 b 各传一个整数,于是按格式”填表”。schema 写得好不好(尤其 description),直接决定模型用得对不对。这也是为什么文件开头把这一层叫「上行翻译层」:把 Python 函数翻译成模型能读懂的 JSON 说明书。

生成的 Tool.parameters 就是一份标准 JSON Schema:

1
2
3
4
5
6
7
8
{
"type": "object",
"properties": {
"a": {"type": "integer"},
"b": {"type": "integer"},
},
"required": ["a", "b"],
}

读作:「参数是一个对象,含 a、b 两个字段,都是整数,都必填」。注意第 56~57 行的循环就是在逐参数填这张「表」。

然后 Tool.to_openai_schema(或 registry.get_schemas 批量)再包一层 OpenAI API 要求的外壳:

1
2
3
4
5
6
7
8
[{
"type": "function",
"function": {
"name": "multiply",
"description": "Multiply two integers.",
"parameters": { ...上面那份 schema... },
},
}]

最终实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
agent.py(循环)      → 只做:调用 registry.execute + 配对写回

registry.py(注册表) → 收完整 tool_call:json.loads + 查表 + 错误转 ToolResult

tools.py(工具本体) → Tool 类:动态建模 + to_openai_schema + execute + __call__

ToolResult(结果) → 永不抛,ok/data/error,serialize 转字符串
---------------------
@tool def get_weather(city: str) -> str: ... # 一行定义

Tool 类(动态建模 + to_openai_schema + execute + __call__)

ToolRegistry.register(tool) / execute(tool_call) # 注册表分发

agent.py: registry.execute(tc).serialize() # 循环只做调用 + 写回
Tool ToolRegistry
视角 单个工具 一群工具
知道其他工具吗 不知道,只管自己 知道全部,管它们的集合
懂协议格式吗 不懂(只收 args: dict) 懂(收完整 tool_call,内部解析)
一句话 「我怎么跑」 「谁在我这里,模型想调谁,我帮它找到并执行」

packages/my-agent-core/src/my_agent_core/tools/core.py

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
ToolResult(dataclass)

@dataclass
class ToolResult:
ok: bool # 是否成功
data: Any = None # 成功时的结果数据
error: str | None = None # 失败时的错误消息
terminate: bool = False # 熔断退出标记(用于阶段 7 提前退出 ReAct 循环)
meta: dict[str, Any] = field(default_factory=dict) # 结构化诊断元数据

def serialize(self) -> str:
# 转成写入 messages 的字符串;失败返回错误文本

Tool(类)

class Tool:
def __init__(
self,
func,
*,
name: str | None = None,
description: str | None = None,
params_model: type[BaseModel] | None = None,
is_parallel_safe: bool = False, # 关键:声明并发安全性(默认 False 保护因果顺序)
timeout: float | None = None, # 执行超时时间(结合 asyncio.wait_for)
raw_schema: dict[str, Any] | None = None, # 外部/MCP 预编译 Schema 透传
): ...

def execute(self, args: dict) -> ToolResult:
# 校验 + 异步执行 + 超时保护,永不抛:pydantic 校验失败或工具异常 → ToolResult(ok=False, error=...)

tool(模块级函数,装饰器工厂)

def tool(
func=None,
*,
name: str | None = None,
description: str | None = None,
params_model: type[BaseModel] | None = None,
is_parallel_safe: bool = False,
timeout: float | None = None,
):
# @tool 装饰器工厂:支持 @tool 和 @tool(name=..., is_parallel_safe=True)
# 内部构造 Tool 对象并返回

my_agent_core/registry.py

1
2
3
4
5
6
7
8
9
10
11
12
ToolRegistry(类)

class ToolRegistry:
def __init__(self): ...

def register(self, tool: Tool) -> None: ... # 注册工具(同名静默覆盖)
def unregister(self, name: str) -> None: ... # 注销工具
def get(self, name: str) -> Tool | None: ... # 查表
def get_schemas(self) -> list[dict]: ... # 导出 OpenAI 兼容 tools 列表

async def execute_batch(self, tool_calls: list[dict]) -> list[ToolResult]:
"""批量执行工具调用:一票否决制因果时序保护调度。"""

核心调度机制:一票否决因果保序并发(Unanimous Parallel)

当大模型在单轮 ReAct 交互中同时发出多个工具调用(例如同时下发 5 个 tool_calls): 1. 全员通过才放行并发: 系统检查该批次目标工具对象的并发标记:all(tool.is_parallel_safe for tool in batch_tools); 只有当该批调用的全部工具均为 is_parallel_safe=True(例如全为只读的 read, grep, find, list),系统才调用 asyncio.gather 并发极速推进,总耗时由 O(N) 骤降至 O(1)! 2. 一票否决回退保序串行: 一旦批次中包含了哪怕一个 is_parallel_safe=False(例如包含写操作 edit, write, todo 或外部未知工具),整批工具调用立刻自动降级为严格按大模型输出的原序串行执行架构价值:从调度器层面从根本上杜绝了“本应先写后读的操作,因并发导致先读到了脏数据”的因果时序倒置(Causal Inversion)风险!

架构演进展望:对标 Tau 的 AgentToolResult 协议升级

在阶段 18 对标 Tau 架构的深度重塑中,工具返回协议将进一步向前迈进: 1. 多模态内容与诊断隔离(AgentToolResult: - content: list[TextContent | ImageContent]:原生向大模型喂回文本块与图片块; - details: JSONValue专门向前端 TUI / CLI 传递诊断元数据(如 exit_code, duration, truncated_bytes),该字段完全不消耗大模型 Token,彻底解耦“给模型看的文本”与“给前端展示的数据”; 2. 动态工具扩展(added_tool_names: 工具(如插件加载器或 MCP 动态发现)执行后,可通过此字段向会话声明动态新开放的工具集,触发上下文窗口的动态 Token 补算; 3. 提前交接终止信号(terminate: 允许工具(如人机交互问卷或不可逆故障)主动向调度循环发出终止信号,无需等待模型多跑一轮。

内部工具

框架层除了让用户自己 @tool,还内置了一批工具——放在 my_agent_core.tools.builtin 包,用工厂函数构造。其中四个文件工具对标 pi 的基本四件套:

工具 签名 作用
read read(path, limit=None) 读文件,limit 按行截断
write write(path, content) 写文件(自动建父目录、覆盖)
edit edit(path, old_text, new_text) 精确替换一处文本
bash bash(command) 在 workspace 根执行 shell

两个共同的关键点

① 工厂函数收 root——路径逃逸防护的边界

四个工具都长这样(工厂函数吃 root,吐 Tool):

1
2
3
4
def make_read_tool(root: str | Path) -> Tool:
def read(path: str, limit: int | None = None) -> str:
...
return Tool(func=read, name="read")

root 是「工作区根目录」,工具内部每条路径都先过一遍 _safe_path

1
2
3
4
5
def _safe_path(root: Path, p: str) -> Path:
path = (root / p).resolve() # 消掉 .. 和符号链接
if not path.is_relative_to(root): # 逃出根目录了?
raise ValueError(f"Path escapes workspace: {p}")
return path

resolve()../etc/passwd 折算成绝对路径,is_relative_to(root) 检查它还在不在根里,逃逸就报错——这是文件工具的第一道安全门。

② 错误不抛,且提供精细化纠错提示(Prompt-Quality Errors)

四个工具内部都是 try/except,把异常转成极具指导意义的错误字符串返回,给模型提供明确的自我纠错线索: - read:越界时明确返回文件实际总行数(如 Offset 200 is beyond end of file ('app.py' has only 80 lines total)); - edit:未找到时提示检查缩进与换行,多处匹配时提示提供更多上下文; - bash:超时时自动捕获并保留超时前打印的已输出日志,方便模型判断是否卡在交互输入(如 -y)。

③ FileMutationQueue 细粒度文件锁(并发安全与性能兼得)

writeedit 工具内部通过 FileMutationQueue 按文件绝对路径获取 asyncio.Lock。修改不同文件时全员并发执行,修改同一文件时自动排队串行,兼备极致性能与写安全性。

bash 的额外两道防护

1
2
3
4
5
6
7
8
_DANGEROUS = ["rm -rf /", "sudo", "shutdown", "reboot", "> /dev/"]

def bash(command: str) -> str:
if any(d in command for d in _DANGEROUS):
return "Error: Dangerous command blocked"
r = subprocess.run(command, shell=True, cwd=root,
capture_output=True, text=True, timeout=120)
return (r.stdout + r.stderr).strip() or "(no output)"
  • 危险命令黑名单rm -rf /sudo 这类直接拦下。是「尽力而为」的字符串匹配,不是真沙箱(真隔离得靠容器/权限层,那是 coding agent 层的事)。
  • 120 秒超时sleep 1000 这种挂死命令不会无限等,超时返回 "Error: Timeout (120s)"

为什么是「工厂函数」而不是「直接一个工具」

因为 root 是显式的、每个应用不同——框架层不知道你的工作区在哪。所以工具是 make_read_tool(root) 现造的,调用方(将来 coding agent 层)传自己的 workspace 根进去,再 Agent(tools=[make_read_tool(root), ...]) 装配。这和「@tool 装饰器直接定义」的区别在于:内置工具需要一个运行时才知道的参数(root)注入。

永不抛出:工具出错也是一条消息

## 一、 观念转变:异常 vs 消息(谁才是接收者?)

理解这一节的关键,在于搞清楚错误是给谁看的:

1
2
3
4
5
6
7
8
❌ 传统思路(把错误当成异常):
工具报错 (FileNotFoundError) ──> 穿透调用栈 ──> 打断 Agent Loop ──> 终端崩溃 ──> 只能由人类程序员重新启动。
【接收者是 Python 解释器 / 调用栈】

✅ Agent 哲学(把错误当成一条消息):
工具报错 (FileNotFoundError) ──> 框架捕获 ──> 包装成普通消息: ToolResultMessage(isError=True)
──> 追加进上下文 ──> 喂给大模型 ──> 模型看懂了:“路径错了,我先 ls 看看目录” ──> 自动修复!
【接收者是大模型】

## 二、 统一出口:6 种错误,1 种产物

在工具调用的完整生命周期中,可能会有 6 个不同阶段抛错,但无论哪一步挂掉,Pi 都保证绝不向外抛异常,全部归一化为一条标准的 ToolResultMessage:

1
2
3
4
5
6
7
8
LLM 输出 ToolCall

├── 1. 工具不存在 (如模型幻觉造了一个未注册的工具) ────> ToolResultMessage { isError: true, content: "Tool xxx not found" }
├── 2. prepareArguments 参数预处理抛错 ───────────────> ToolResultMessage { isError: true, content: 预处理异常 }
├── 3. Schema 参数强校验失败 (如传错类型) ───────────> ToolResultMessage { isError: true, content: Pydantic 校验错误 }
├── 4. beforeToolCall 权限拦截 (如危险命令被阻断) ─────> ToolResultMessage { isError: true, content: 拦截原因 }
├── 5. tool.execute 运行时崩溃 (如 500/超时/文件不存在) ─> ToolResultMessage { isError: true, content: 运行时异常 }
└── 6. afterToolCall 后置处理抛错 ────────────────────> ToolResultMessage { isError: true, content: 后置异常 }

对 Agent Loop 而言:它看到的结果永远是一个干净的 ToolResultMessage,因此主循环的 while 可以安全地继续转动,把结果喂给下一轮大模型。

经典对比(看看 Pi 是怎么写第一层的):

### 1. Read 工具(read.ts)

  • ❌ 差的报错:raise Exception(“Read error”) → 模型两眼一抹黑。
  • ✅ Pi 的报错:Offset 200 is beyond end of file (100 lines total) → 模型立即明白:“文件只有 100 行,那我下次传 offset=50”。

### 2. Bash 工具(bash.ts)

  • ❌ 差的报错:raise Exception(“Command failed”)。
  • ✅ Pi 的做法(教科书级): 把“执行到一半已经被捕获的 stdout 输出” + “退出码 / 超时状态” 打包在一起返回:
    1
    2
    3
    4
    5
    Command failed with exit code 1.
    --- Output before error ---
    npm ERR! code ENOENT
    npm ERR! syscall open
    npm ERR! path /package.json
    模型看到这段具体的错误输出,就能像人一样分析根因。

异常 → 消息:编码前后对比

1
2
3
4
5
6
7
8
9
10
11
12
13
14
工具抛出的原始异常(except 之前):       编码后的 ToolResultMessage(except 之后):
FileNotFoundError: [Errno 2] No such {
→ 一路穿透管道 role: "toolResult",
→ 打断 Agent Loop toolCallId: "call_abc",
→ 事件序列不完整,UI 卡死 toolName: "read",
content: [{
type: "text",
text: "[Errno 2] No such file or directory"
}],
isError: True ← 唯一标记
}
→ 追加到对话历史
→ 下一轮发给模型
→ 模型看到后自己决定怎么办

写自定义工具时的最佳实践

借鉴 Bash 工具的写法,自定义工具的 execute 应该长这样:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
async def execute_custom_tool(params: MyArgs) -> ToolResult:
try:
# 1. 正常业务逻辑
result_data = await do_some_io(params)
return ToolResult(ok=True, data=result_data)

except KnownBusinessErrorA as err:
# 2. 主动识别错误 A:给出具体的线索与修复建议
return ToolResult(
ok=False,
error=f"Resource '{params.id}' not found. Available IDs are: {get_available_ids()}"
)

except KnownBusinessErrorB as err:
# 3. 主动识别错误 B:给出范围限制
return ToolResult(
ok=False,
error=f"Timeout after {params.timeout}s. Try increasing timeout or splitting query."
)

except Exception as err:
# 4. 未知异常:尽量带上原始异常细节,交给框架
return ToolResult(ok=False, error=f"Unexpected failure: {type(err).__name__}: {err}")


终极防线:工具尚未执行就被中途中断?——对标 Tau 的 tool_history 三阶段自愈状态机

在前一章中,我们详细分析了“永不抛出(Never-Throw)”原则——它保证了当代码执行进工具函数内部时,即便遇到异常也会被包装为 ToolResult(ok=False, error=...),避免 Agent 进程崩溃。

但是在真实大模型 Agent 系统中,还潜伏着一个更隐蔽、破坏力更强的系统级死穴:如果工具压根“没来得及执行完”,或者在执行前夕就被外部打断了呢?

1. 致命的“断头工具调用(Dangling Tool Call)”

各大主流大模型(OpenAI、Anthropic、DeepSeek 等)在设计 Function Calling 协议时,有着严苛的图灵机契约:

🚨 模型 API 契约约束: 如果一条 assistant 消息声明了 tool_calls: [{"id": "call_123", ...}],那么在对话历史中,其紧邻的下一条消息必须是 role: "tool" 且带上匹配的 tool_call_id: "call_123"

如果出现以下任何一种真实突发场景: 1. 用户主动取消:模型发起了一个耗时 30 秒的编译命令工具,用户等不及直接按了 Ctrl+C 或触发了 await agent.abort(); 2. 网络异常截断:模型吐出了工具调用描述后,网络突然抖动中断; 3. 并发工具部分失败:模型一次性发起了 3 个并发工具调用,第 1 个执行成功,第 2 个发生致命系统错误直接退出了当轮 ReAct 迭代。

这会导致会话历史里留下了一条“只有 ToolCall,没有对应 ToolResult”的消息——即断头调用

现实案例剖析:一个“断头”如何让会话彻底脑死亡?

我们来看一个真实发生过的典型崩溃场景:

  1. 模型决定调用工具:一次性生成了两个工单 call_001(读 a.txt)和 call_002(一个耗时 30 秒的巨型编译命令)。
  2. 正常执行前半截call_001 顺利执行完成,结果回填。
  3. 中途意外打断:正在执行 call_002 时,用户等不及了按下了 Ctrl+C,或者代码里触发了 await agent.abort()
  4. 致命空档产生:此时 Python 进程被打断,call_002 根本没有产生任何返回值
  5. 坏死数据固化:会话历史直接停留在这一秒,并持久化到了硬盘的 .jsonl 文件中。

此时磁盘里的消息记录变成了这种残缺状态(断头了!):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
[
{"role": "user", "content": "帮我处理一下代码"},
{
"role": "assistant",
"tool_calls": [
{"id": "call_001", "name": "read"},
{"id": "call_002", "name": "build"} // ← 发起了 call_002
]
},
{
"role": "tool",
"tool_call_id": "call_001",
"content": "a.txt 的内容是 Hello"
}
// 🚨 致命缺失:由于中途中断,历史里根本没有 call_002 的结果!
]

接下来会发生什么悲剧?

当下一次用户再发一句新的话(比如:“算了,不要编译了,给我讲个笑话吧”): Agent 会把上面这段包含断头的历史加上新的提问,一起打包发给 OpenAI / DeepSeek / Claude。

大模型服务端的网关校验器一扫历史,发现 assistant 发起了 call_002,后面居然没有对应的 tool 结果,网关直接无情抛出:

1
2
3
HTTP 400 Bad Request:
An assistant message with 'tool_calls' must be followed by tool messages responding to each 'tool_call_id'.
Missing tool response for: call_002.

更致命的是:因为这条残缺的历史记录已经保存在你的硬盘会话文件里了,以后只要你加载这个会话,无论发什么新指令,大模型永远报 400!这个会话文件就彻底“脑死亡”报废了!

tool_history.py 是如何化解这场悲剧的?

tool_history.py 就像一个“消息转录本的智能外科医生兼安检门”。在会话从磁盘恢复后、以及每次送给大模型之前,它都会对整个消息链进行拓扑自愈。

当它扫描到上述残缺历史时,发现 call_002 悬空无下文,它会在内存中自动就地合成一条合法的工具结果插进去:

1
2
3
4
5
6
{
"role": "tool",
"tool_call_id": "call_002",
"content": "Tool call interrupted by user", // 明确告诉大模型:这个工具被用户打断了
"metadata": {"is_error": true}
}

这样一来,两全其美: 1. 大模型的 API 契约瞬间被满足了:每个 call 都有对应的 tool 结果紧随其后,API 绝对不会再报 400 拒绝服务! 2. 大模型的认知逻辑也顺畅了:模型看到 Tool call interrupted by user,在上下文中就自然理解“哦,原来刚才那个任务被用户取消了”,下一轮它就能基于这个事实正常回答!


2. 为什么简单的遍历补齐搞不定?

最朴素的想法是:“遍历消息列表,只要发现某条 Assistant 消息有 tool_calls,后面没跟 Tool 消息就硬塞一条假消息进去。”

在严苛的工程实践中,这种写法会踩进大坑: 1. 跨轮同名 ID 复用:某些开源或商用模型在多轮长对话中,可能会重复生成相同的 tool_call_id(如多次调用都叫 "call_0")。朴素遍历极易把第 5 轮的真实合法结果“偷”给第 1 轮的同名调用,导致第 5 轮反而变成了断头! 2. 位置错位与孤儿结果(Orphan Results):偶尔因并发乱序或重试,存在没有被任何 Assistant 声明引用的游离 ToolResult。如果在上下文中放行,同样会触发 API 400 报错。


3. 对标 Tau 的三阶段确定性拓扑自愈算法

为了彻底扫清断头死锁,我们在阶段 17 深度对齐了 Tau (tau_agent.tool_history) 的三阶段确定性自愈状态机,并在核心层实现了 my_agent_core/tool_history.py

算法核心流转拓扑如下:

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
原始可能有断头/孤儿的历史消息列表


┌─────────────────────────────────────────────────────────────┐
│ 【Phase 1: 预留就近配对 (Reserve Adjacent Pairs)】 │
│ 先扫描当前位置与期望位置完全匹配的合法结果,建立强绑定排他 │
│ 锁定,占住位点,绝对防止被跨轮同名 ID 或贪心扫描错误抢夺! │
└──────────────────────────────┬──────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 【Phase 2: 贪心匹配与断头补齐 (Greedy Match or Synthesize)】 │
│ 对所有未被锁定的 ToolCall 进行向后贪心扫描: │
│ • 优先在调用位置之后寻找合法的真实结果; │
│ • 若彻底找不到,自动合成标准中断结果: │
│ Message(role="tool", content="Tool call interrupted by │
│ user", metadata={"tool_call_id": id, │
│ "is_error": True}) │
└──────────────────────────────┬──────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 【Phase 2.5: 真实结果反超 (Real Result Priority)】 │
│ 若后续扫描发现由于遍历顺序遗漏了真实合规的结果,立即撤销 │
│ 已合成的中断占位,让真实结果反超生效! │
└──────────────────────────────┬──────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 【Phase 3: 转录本重塑与孤儿清理 (Reconstruction & Pruning)】 │
│ • 按照 Assistant 声明的调用顺序,严格紧随插入匹配的工具结果 │
│ • 凡是没有被任何 ToolCall 认领的游离结果(孤儿),全部剔除!│
└─────────────────────────────────────────────────────────────┘


100% 结构合法的自愈转录本

4. 核心源码落地与诊断模型 (tool_history.py)

自愈函数返回一个结构化不可变诊断数据类 ToolHistoryRepair

1
2
3
4
5
6
7
8
9
10
11
12
13
# packages/my-agent-core/src/my_agent_core/tool_history.py

_INTERRUPTED_TOOL_RESULT = "Tool call interrupted by user"

@dataclass(frozen=True, slots=True)
class ToolHistoryRepair:
"""修复后的合法转录本以及结构化诊断计数。"""
messages: tuple[Message, ...] # 自愈后的合法消息序列
changed: bool = False # 是否发生了纠偏与修复
synthesized_results: int = 0 # 补齐的中断结果条数
dropped_orphan_results: int = 0 # 丢弃的游离孤儿结果条数
dropped_duplicate_results: int = 0 # 丢弃的重复结果条数
reordered_results: int = 0 # 纠正乱序重排的条数

核心修复算法 repair_tool_history 的精简实现:

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
def repair_tool_history(messages: Sequence[Message]) -> ToolHistoryRepair:
"""对会话历史进行确定性拓扑自愈,确保所有工具调用均合法闭合。"""
# 提取所有 (msg_idx, call_offset) 的工具调用事件
call_occurrences = ...
# 按 tool_call_id 索引所有真实存在的 tool 消息位置
results_by_id = ...

selected_results: dict[tuple[int, int], tuple[int | None, Message]] = {}
used_result_positions: set[int] = set()
synthesized_results = 0

# ── Phase 1: 预留已就近配对的调用 ──
for occurrence, call, expected_pos in call_occurrences:
if expected_pos < len(messages) and _get_tool_call_id(messages[expected_pos]) == call.get("id"):
selected_results[occurrence] = (expected_pos, messages[expected_pos])
used_result_positions.add(expected_pos)

# ── Phase 2: 剩余调用贪心匹配或补齐中断结果 ──
for occurrence, call, _ in call_occurrences:
if occurrence in selected_results:
continue
call_id = str(call.get("id", ""))
candidates = results_by_id.get(call_id, [])
matched_pos, matched_msg = _find_best_candidate(candidates, used_result_positions, occurrence[0])

if matched_msg is not None:
selected_results[occurrence] = (matched_pos, matched_msg)
used_result_positions.add(matched_pos)
else:
# 补齐标准中断工具结果
synthetic = Message(
role="tool",
content=_INTERRUPTED_TOOL_RESULT,
metadata={"tool_call_id": call_id, "is_error": True},
)
selected_results[occurrence] = (None, synthetic)
synthesized_results += 1

# ── Phase 2.5: 真实结果反超合成中断 ──
# ... 确保真实结果绝不被虚假中断误吞 ...

# ── Phase 3: 重建转录本、孤儿丢弃与保序重排 ──
repaired: list[Message] = []
for msg_idx, message in enumerate(messages):
if message.role == "tool":
# 游离孤儿丢弃;已配对工具在 assistant 之后紧邻插入,此处直接跳过
continue

repaired.append(message)
if message.role == "assistant":
for offset in range(1, len(_get_tool_calls(message)) + 1):
_, tool_res = selected_results[(msg_idx, offset)]
repaired.append(tool_res) # 严格保序紧随其后

return ToolHistoryRepair(messages=tuple(repaired), changed=..., synthesized_results=synthesized_results, ...)

5. 架构级双重防护网:内层 Never-Throw 护函数,外层 tool_history 护转录本

通过这次 Tau 对齐演进,整个工具调用系统形成了固若金汤的双层立体防御架构

防御层级 核心护卫模块 拦截时序 防护目标 最终结果
内层执行防御 Tool.execute() (Pydantic + Try/Catch) 工具函数执行中 参数非法、业务抛错、OS 路径逃逸、超时 转为 ToolResult(ok=False) 反馈大模型自纠,Agent 循环不崩
外层拓扑防御 tool_history.py (repair_tool_history) 会话反序列化 / agent.abort() / 调用 LLM 前夕 用户按 Ctrl+C 中断、网络断流、断头调用、孤儿结果 确定性自愈重构转录本,大模型 API 永不报 400

至此,哪怕用户在工具并发执行的任意毫秒暴力掐断任务,或者强行杀死进程后重启系统,会话也能在毫秒级内自动恢复闭合,彻底攻克了大模型应用中最棘手的会话持久化死锁顽疾!


6. 工业级七阶段工具执行流水线:并发批处理与实时进度流

在完成了工具参数校验和拓扑自愈后,当大模型在一轮中发起了多个工具调用时,我们进入了最核心的 工具执行车间 (_execute_tools_turn)

严格对齐 Pi 官方架构契约与 Tau 微内核设计,工具执行被划分为七个确定性阶段:

  1. 阶段 1:输出截断防御检查 (_fail_tool_calls_from_truncated_message): 当检测到模型输出触达 Token 上限被截断(stop_reason == "length")时,断然拒绝执行任何工具,防止流式 salvage 拼出残缺参数导致代码写崩或命令腰斩。自动合成警告错误并回传模型引导重新完整发起调用。
  2. 阶段 2:Preflight 广播 (ToolExecutionStart): 在审批与执行前,率先按 source order 广播 ToolExecutionStart,使 UI 能够毫秒级渲染工具准备运行状态。
  3. 阶段 3:串并行决策网关 (ToolRegistry.execute_batch): 悲观读写分流:全只读安全工具(is_parallel_safe=True)启用 asyncio.gather 全并发加速;只要包含任一写入/串行工具,整批退化为保序串行执行,防止因果时序倒置。
  4. 阶段 4:前置审查审批与改参 (before_tool_call): 通过 _coerce_tool_call 归一化入参,调用 before_tool_call 审批,支持安全阻断(block)与参数就地热修改(updated_args)。
  5. 阶段 5:并发批执行与流式进度回传 (ToolExecutionUpdate): 采用 asyncio.Queueloop.call_soon_threadsafe 跨线程安全桥接,支持长耗时工具(如 Bash 编译或子代理)在运行态向外广播累积快照(Cumulative Snapshot)。生命周期锁存(accepting_updates)确保工具返回后丢弃迟到回调。
  6. 阶段 6:后置改写与单工具终态广播 (after_tool_call & ToolExecutionEnd): 调用 after_tool_call 支持结果脱敏与改写(updated_result),广播包含 terminate 状态的 ToolExecutionEnd
  7. 阶段 7:转录本保序归档与批量优雅熔断 (MessageStart/End & should_terminate): 无论并发执行完成顺序如何,回传大模型的 role="tool" 消息严格按 Assistant 原始 Source Order 恢复排布。若批次中任一工具(any() 语义)或 Hook 返回 terminate=True,立即终结 ReAct 循环,保全 final_text 并正常交付结果。

阶段 5 核心攻坚:生产者-消费者管道与三大技术死结破解

在实现第 5 阶段(并发执行与实时流式进度回传)时,架构面临了三个看似不可调和的技术冲突:

1
2
3
4
5
6
7
8
9
10
11
12
13
【冲突 1】异步生成器 vs 并发批量执行
• 需求 A:多个工具必须并发跑(比如同时读 3 个文件),底层用的是 asyncio.gather(...);
• 需求 B:_execute_tools_turn 是一个生成器,必须一有进度就立刻 yield 出去。
• 痛点:在 asyncio.gather 的深处是不能直接向外层生成器 yield 的!

【冲突 2】主事件循环 vs 工作线程池
• 需求 A:很多工具是同步的(如普通的 Python 函数、调用操作系统的 bash),必须扔进线程池
asyncio.to_thread 跑,不能卡死主线程;
• 痛点:Python 的异步队列 asyncio.Queue 是严格单线程的,工作线程直接碰队列就会崩溃闪退!

【冲突 3】实时流式读取 vs 任务结束防死锁
• 需求 A:前台必须一直监听队列,只要有日志就拿出来;
• 痛点:前台怎么知道后台什么时候“全部跑完了”?如果盲目等,后台跑完后前台就会永久卡死(死锁)。

破局方案:前后台解耦的“单向传送带”管道模型

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
┌──【后台 Producer: runner 异步任务】──┐
│ │
│ 并发执行 tool_1, tool_2, tool_3 ... │
│ │ │
│ ▼ 产生进度 │
│ on_update("50% complete") │
│ │ │
│ ▼ │
│ safe_put_update (线程安全路由) │
│ ├─ 主线程: queue.put_nowait() │
│ └─ 工作线程: call_soon_threadsafe │
│ │ │
│ ▼ 写入 │
└─────────────> ┌────────┐ <────────────┘
│ 异步队列│
│ queue │
└────────┘

▼ 读取
┌──【前台 Consumer: 生成器主循环】────┐
│ │
│ while True: │
│ item = await queue.get() │
│ if item is _SENTINEL: break │
│ yield item (实时推给前端 UI!) │
│ │
└───────────────────────────────────────┘

核心代码落地:5 个严丝合缝的实现步骤

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
# ── 阶段 B: 并发批执行与实时进度流式广播 ──
if prepared_calls:
# 步骤 1: 建立线程安全事件队列与独一无二的哨兵
queue: asyncio.Queue[Event | object] = asyncio.Queue()
_SENTINEL = object()
loop = asyncio.get_running_loop()
loop_thread_id = threading.get_ident()

# 步骤 2: 跨线程安全网关路由
def safe_put_update(ev: Event) -> None:
if threading.get_ident() == loop_thread_id:
queue.put_nowait(ev)
else:
with contextlib.suppress(RuntimeError):
loop.call_soon_threadsafe(queue.put_nowait, ev)

# 步骤 3: 柯里化闭包工厂,给每个工具定制专属 on_update
def make_on_update(call_id: str, tool_name: str, args: dict[str, Any]) -> Callable[[Any], None]:
def on_update(partial: Any) -> None:
safe_put_update(
ToolExecutionUpdate(
tool_call_id=call_id,
tool_name=tool_name,
args=args,
partial_result=partial,
)
)
return on_update

calls_to_run = [
(call.name, current_args, make_on_update(call.id, call.name, current_args), call.id)
for _, call, current_args in prepared_calls
]

# 步骤 4: 后台并发批处理任务(无论成败,finally 必送哨兵防死锁)
async def _run_batch() -> list[ToolResult]:
try:
return await registry.execute_batch(calls_to_run, signal=signal)
finally:
if threading.get_ident() == loop_thread_id:
queue.put_nowait(_SENTINEL)
else:
with contextlib.suppress(RuntimeError):
loop.call_soon_threadsafe(queue.put_nowait, _SENTINEL)

runner = asyncio.create_task(_run_batch())
try:
# 前台消费循环:见事件就 yield,见哨兵就 break
while True:
item = await queue.get()
if item is _SENTINEL:
break
if isinstance(item, Event):
yield item
batch_out = await runner
for (idx, _, _), res in zip(prepared_calls, batch_out, strict=False):
direct_results[idx] = res
except Exception as exc:
for idx, _, _ in prepared_calls:
if idx not in direct_results:
direct_results[idx] = ToolResult(ok=False, error=f"Tool execution failed: {exc}")
finally:
# 步骤 5: 外层中断时的防暴清场
if not runner.done():
runner.cancel()
with contextlib.suppress(asyncio.CancelledError, Exception):
await runner

这套架构将“不可流式的多任务并发”优雅转化为“可流式的单通道事件流”,用 call_soon_threadsafe 抹平了同步线程与异步主循环的鸿沟,用 _SENTINEL 彻底封死了死锁漏洞,达到了真正的工业级水准!


深度辨析:宏观串并行 vs 微观同异步(及异步队列的本质)

很多开发者在理解工具并发与流式回传时,容易把“宏观调度的串并行”“底层函数的同异步”混淆。实际上,这是两个完全正交的独立维度:

维度 关心的核心问题 决定者 核心手段
维度 1:宏观批处理
(并发还是串行?)
这批工具是一起开工,还是一个接一个排队? ToolRegistry.execute_batch
(基于 is_parallel_safe 属性)
• 并发:asyncio.gather(...)
• 串行:for 循环依次 await
维度 2:单工具执行
(主线程还是子线程?)
这个工具自身的 Python 代码会卡死主线程吗? Tool.execute
(基于 inspect.iscoroutinefunction
• 异步(async def):主线程协程直接跑
• 同步(普通 def):扔进系统线程池 asyncio.to_thread

1. 核心澄清:异步队列里流动的到底是什么?

请务必注意: 异步队列 queue 里面装的不是工具执行的最终返回值(ToolResult
最终结果是在所有工具都跑完后,由 batch_out = await runner 一次性拿到的。
队列里流动的,纯粹是工具在运行中途发射出来的“中间流式进度事件”(ToolExecutionUpdate

2. 如果工具是【并行】,底层具体怎么运转?

假设模型单轮发起了 3 个工具调用:tool_1async def 协程)、tool_2(普通 def 计算)、tool_3(普通 def 命令),且均声明为 is_parallel_safe=True: - 宏观调度registry.execute_batch 使用 asyncio.gather 将它们同时打出去并发运行; - 微观分流tool_1 在主事件循环中跑协程;tool_2tool_3asyncio.to_thread 分配给系统线程池中的两个独立子线程并行跑; - 进度汇入: - tool_1 在主线程调用 on_update safe_put_update 识别为主线程 queue.put_nowait 直接丢入; - tool_2 在子线程调用 on_update safe_put_update 识别为工作线程 通过 loop.call_soon_threadsafe 安全跨线程预约投递; - 前台感知:前台 while True: item = await queue.get() 无论谁的进度先到,就立刻把谁先 yield 出来给外部界面,呈现出多个任务交织滚动的极致流式体验。

3. 如果工具是【串行】,又是如何处理的?

假设模型发起了一个写操作 edit_fileis_parallel_safe=False)和一个读操作 read_file: - 宏观调度:触发“一票否决”,registry.execute_batch 退化为顺序遍历:for tc in tool_calls: await self.execute_tool(tc, ...); - 微观分流与队列复用: - 先执行 edit_file:中途产生的进度事件依次塞入 queue 前台立刻实时打印,直到 edit_file 彻底结束; - 接着执行 read_file:中途产生的进度事件塞入同一个 queue 前台接着打印,直到 read_file 结束; - 整批串行任务结束:触发 finally 塞入 _SENTINEL 哨兵,前台收工退出; - 整套传送带与哨兵机制 100% 无缝复用,零多余逻辑


实战指南:什么样的工具写成 async def?什么样的工具写成普通 def?

在 Agent 系统的工程实践中,判断工具应该写成异步(async def还是同步(普通 def,核心依据只有一个:

这个工具在执行时,主要是在“等外部事件(网络/其他服务/子代理)”,还是在“让本地 CPU/操作系统猛跑”?

1. 必须 / 适合写成【异步工具】(async def

核心特征:需要长时间等待外部世界响应,等待期间可以主动让出 CPU 控制权(await),让主线程去处理其他任务。

  • 子代理委派(Subagent / Task):例如 task(prompt="审查代码", agent="reviewer")。子 Agent 还要经历自己的推理与工具循环,耗时可达数秒甚至数分钟,必须通过 await subagent.run() 异步挂起。
  • 网络检索与外部 API(Web Search / Crawl / MCP 协议):底层依赖 httpx.AsyncClientaiohttp,等待远程服务器网络握手与数据传输。
  • 浏览器自动化(Playwright / Browser):等待页面加载(domcontentloaded)、等待元素渲染、网络空闲等,天然全是非阻塞异步调用。
  • 定时器与轮询工具:内部使用 await asyncio.sleep(delay) 优雅挂起,绝不卡死主线程。

2. 通常写成【同步工具】(普通 def

核心特征:调用的几乎全是本地现成计算资源、操作系统 API 或传统三方阻塞库,代码自上而下顺次执行。

  • 本地快速小文件读写与编辑(read / write / edit:在现代 SSD 上读写几十 KB 文本只需 0.1~0.5 毫秒,直接使用 open()pathlib.Path.read_text() 简单直观,无需额外创建异步事件调度开销。
  • CPU 密集型计算与正则扫描(grep / calculator / ast-grep:底层是纯 CPU 满负荷运算或 C/Rust 原生扩展,中途根本没有空闲等待时间,写成 async def 毫无意义(无处可 await)。
  • 简单封装传统阻塞库的系统工具(bash:许多开发者习惯直接调用 subprocess.run(cmd, shell=True),这种阻塞操作适合作为普通 def

3. 框架的“零心智负担”自适应无感桥接

为了让工具开发者彻底摆脱心智负担,my-pi-agentTool.execute 内部建立了自适应分流机制:

1
2
3
4
5
6
if self.is_async:
# async def: 直接在主事件循环中高性能协程 await
result = await async_func_call()
else:
# 普通 def: 框架自动封装进 asyncio.to_thread,发配给后台线程池执行,绝对不阻塞主事件循环!
result = await asyncio.to_thread(func_call)

工程收益:开发者想怎么写就怎么写。写同步无需担心卡死 Agent,写异步能够极致榨干事件循环吞吐量。

场景特点 建议定义为 典型例子 框架底层的实际执行环境
需要联网 / 调远程服务 async def web_searchhttp_fetchmcp_client 主事件循环(协程非阻塞挂起)
需要调用子 Agent / 嵌套 Agent async def taskdelegate_subagent 主事件循环(协程非阻塞挂起)
浏览器自动化 async def agent_browserplaywright_click 主事件循环(协程非阻塞挂起)
本地读写文件 / 简单计算 普通 def readwriteeditcalculate 线程池 asyncio.to_thread(极速跑完)
运行 Shell / 本地进程 普通 def bashsubprocess.run 线程池 asyncio.to_thread(后台运行)
纯正则搜索 / AST 语法分析 普通 def grepast_grep_search 线程池 asyncio.to_thread(算完返回)

7. 深入底层:并发调度与同异步执行全景(宏观批处理 vs 单工具执行)

在阅读工具执行源码时,很多人容易把“宏观调度的串并行”“底层函数的同异步”混为一谈,甚至产生误解:“是不是并行工具就是用子线程跑,串行工具就是直接跑?”

实际上,在 my-pi-agent 的工具流水线中,这是两个完全解耦的独立正交维度

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
┌─────────────────────────────────────────────────────────────────────────────┐
│ 两个完全正交的架构维度 │
├─────────────────────────────────────────────────────────────────────────────┤
│ │
│ 【维度 1:宏观批处理 (Macro Batch Scheduling)】 │
│ • 关心的核心问题:这批工具是一起开工,还是一个接一个排队? │
│ • 决定者与位置:ToolRegistry.execute_batch (基于 is_parallel_safe 属性) │
│ • 实现手段: │
│ - 全并发:整批全为 is_parallel_safe=True ➔ asyncio.gather(...) 并发 │
│ - 保序串行:含任一写入工具 ➔ 悲观回退,for 循环依次 await │
│ │
│ 【维度 2:单工具内部执行 (Micro Single-Tool Execution)】 │
│ • 关心的核心问题:这个工具自身的 Python 代码会卡死事件循环主线程吗? │
│ • 决定者与位置:Tool.execute (基于 inspect.iscoroutinefunction 静态自省) │
│ • 实现手段: │
│ - 原生异步:async def ➔ 直接 await tool.func(...) (主线程协程调度) │
│ - 同步阻塞:普通 def ➔ asyncio.to_thread(_sync_run) 扔进线程池 │
│ │
└─────────────────────────────────────────────────────────────────────────────┘

1. 维度 1:宏观批处理调度(并发还是串行?)

ToolRegistry.execute_batch 全权掌舵。其采用“一票否决制”的悲观读写分流策略

  1. 全并发加速(asyncio.gather
    • 触发条件:大模型在一轮中呼叫的所有工具,其 tool.is_parallel_safe 均为 True(例如同时读取 3 个只读文件、查询 2 个外部只读 API)。
    • 调度方式:使用 asyncio.gather(*[self.execute_tool(...) for ...]) 一起并发调度,整体耗时从累加缩减为“取决于最慢的那一个”。
  2. 保序串行(悲观回退)
    • 触发条件:这一批调用中哪怕包含任一一个 is_parallel_safe=False 的工具(如写入文件 write、编辑文件 edit、运行终端命令 bash)。
    • 调度方式:整批工具立刻放弃并发,严格退化为按大模型在提示词里的原始声明顺序(Source Order),使用普通的 for 循环逐个串行执行
    • 架构不变式:严防因果倒置(例如大模型本意是“先编辑代码,再执行测试”,如果盲目并发,可能导致测试在代码还没写完前就抢跑报错)。

2. 维度 2:单工具内部执行(异步协程还是线程池?)

无论宏观上是 asyncio.gather 并发还是 for 循环串行,每个具体工具在执行时,依然由其自身的实现形态决定在哪个线程运行

1
2
3
4
5
6
7
# Tool.execute 的核心分流逻辑
if inspect.iscoroutinefunction(self.func):
# 原生异步工具:零线程开销,直接在当前主事件循环上跑
result = await self.func(...)
else:
# 同步阻塞工具:坚决不能卡死主线程,封装进线程池独立子线程跑
result = await asyncio.to_thread(_sync_run)

怎么判断工具该写成【异步工具(async def)】还是【同步工具(普通 def)】?

核心依据只有一个:这个工具在执行时,主要是在“等外部世界响应”,还是在“占用本地 CPU 或操作系统的同步进程”?

工具类型 函数签名 典型应用场景 内部执行机制 为什么这么写?
异步工具 async def • 子代理委派(task
• 网络请求(fetch_web / aiohttp)
• 异步数据库驱动(asyncpg)
直接在主事件循环中 await,不占用额外工作线程 子代理和网络请求动辄等待几秒至几十秒,期间主动让出 CPU 控制权,主线程可并发处理其他事件广播或打断信号。
同步工具 普通 def • 本地文件操作(read / write / edit
• 纯 CPU 计算与文本处理(math / json)
• 本地 Shell 阻塞子进程(bash subprocess)
框架自动通过 asyncio.to_thread 投入底层线程池 Python 的 subprocess.run 或标准 open() 是同步阻塞的系统调用,如果不扔进子线程,主线程事件循环会被直接焊死,导致打字机流式输出和 UI 动画瞬间冻结!

3. 两套机制交织碰撞:异步队列与跨线程安全

现在把两张拼图合在一起,整个底层流式通信的全景就清晰无比了:

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
┌─────────────────────────────────────────────────────────────────────────────┐
│ 宏观调度 (execute_batch): 决定是 asyncio.gather 还是串行 for 循环 │
└──────────────────────────────────────┬──────────────────────────────────────┘
│ 遍历调起单个工具 (execute_tool)

┌─────────────────────────────────────────────────────────────────────────────┐
│ 单工具执行 (Tool.execute): │
│ ├─ Case A: 若工具是 async def (如 task) ────────► 在主线程跑 │
│ │ │ │
│ │ ▼ 调 on_update │
│ │ 直接 queue.put_nowait() │
│ │ │
│ └─ Case B: 若工具是同步 def (如 bash) ────────► 扔进工作子线程跑 │
│ │ │
│ ▼ 调 on_update │
│ 跨线程必须用 call_soon_threadsafe 桥接! │
└──────────────────────────────────────┬──────────────────────────────────────┘
│ 全部安全汇入

┌──────────────────────────────────┐
│ 主线程专属 asyncio.Queue (传送带) │
└─────────────────┬────────────────┘
│ await queue.get()

┌─────────────────────────────────────────────────────────────────────────────┐
│ 外部异步生成器 (_execute_tools_turn): 实时 yield ToolExecutionUpdate 给 UI │
└─────────────────────────────────────────────────────────────────────────────┘
  • 无论这批工具在宏观上是并发还是串行:后台都运行在独立的 runner = asyncio.create_task(...) 任务中,前台始终是 while True: item = await queue.get() 的流畅消费者;
  • 无论具体工具是运行在主线程的协程、还是运行在子线程里的阻塞代码safe_put_update 都能通过 threading.get_ident() == loop_thread_id 自动识别,丝滑抹平线程边界,保证事件安全、保序地呈现在用户眼前!

pi架构详解

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
pi 启动

├─► project_trust(仅用户/全局和 CLI 扩展生效,在加载项目资源之前触发)
├─► session_start { 原因: "startup"(启动) }
└─► resources_discover { 原因: "startup"(启动) }


用户发送 prompt ───────────────────────────────────────────┐
│ │
├─►(优先检查扩展命令,若匹配则直接旁路/绕过) │
├─► input(输入事件:可拦截、转换或直接处理) │
├─►(若未被处理,进行 Skill/Template 技能与模板展开) │
├─► before_agent_start(可注入消息、修改系统提示词) │
├─► agent_start │
├─► message_start / message_update / message_end │
│ │
│ ┌─── turn 轮次(在 LLM 调用工具期间重复执行) ───┐ │
│ │ │ │
│ ├─► turn_start │ │
│ ├─► context(可修改消息列表/上下文) │ │
│ ├─► before_provider_headers(可修改请求头) │ │
│ ├─► before_provider_request(可检查或替换请求载荷 Payload)
│ ├─► after_provider_response(响应状态码 + 响应头,在消费流之前)
│ │ │ │
│ │ LLM 给出响应,可能会调用工具: │ │
│ │ ├─► tool_execution_start │ │
│ │ ├─► tool_call(可进行阻断/拦截) │ │
│ │ ├─► tool_execution_update │ │
│ │ ├─► tool_result(可修改工具返回结果) │ │
│ │ └─► tool_execution_end │ │
│ │ │ │
│ └─► turn_end │ │
│ │
├─► agent_end │
└─► agent_settled(无剩余重试/压缩/排队追加消息,进入沉淀等待状态)

用户发送下一条 prompt ◄────────────────────────────────────┘

/new(新建会话)或 /resume(切换会话)
├─► session_before_switch(可取消切换)
├─► session_shutdown(关闭旧会话)
├─► session_start { 原因: "new" | "resume", 上一次的会话文件? }
└─► resources_discover { 原因: "startup" }

/fork 或 /clone(分叉或克隆会话)
├─► session_before_fork(可取消分叉)
├─► session_shutdown
├─► session_start { 原因: "fork", 上一次的会话文件 }
└─► resources_discover { 原因: "startup" }

/name 或 pi.setSessionName()(会话重命名)
└─► session_info_changed(会话信息变更)

/compact 或自动压缩(上下文压缩)
├─► session_before_compact(可取消或自定义压缩逻辑)
└─► session_compact

/tree 历史树导航
├─► session_before_tree(可取消或自定义导航逻辑)
└─► session_tree

/model 或 Ctrl+P(选择/循环切换模型)
├─► thinking_level_select(若切换模型导致思考深度发生变化或被截断)
└─► model_select

思考深度变更(通过设置、快捷键绑定或 pi.setThinkingLevel() 触发)
└─► thinking_level_select

退出(Ctrl+C、Ctrl+D、SIGHUP、SIGTERM)
└─► session_shutdown(会话关闭清理)

pi的这个设计和langchain的Middleware十分类似

如果把这些控制逻辑全写在 Agent 的主循环里,代码会变得极其臃肿。因此,Pi 的 Extension 系统与 LangChain Agent Middleware 都选择将控制权(Control Layer)与执行层(Execution Layer)解耦,在 Agent 生命周期的关键切面上暴露钩子。

image-20260801154819449
功能维度 LangChain Middleware 机制 Pi Extension 生命钩子 (Hooks) 共同解决的场景
Agent 启动入口 before_agent input, before_agent_start 拦截用户输入、预处理 Prompt
上下文修改 before_model, modify_model_request context, before_provider_request 上下文裁剪(Compaction)、RAG 动态注入、Payload 替换
网络/Header 控制 Client Transport Interceptor before_provider_headers, after_provider_response 动态切换 API Key、注入自定义 Header、监听响应头
工具调用拦截 after_model / HumanInTheLoopMiddleware tool_call (可直接返回 block) 人工审批(HITL)、安全阻断拦截
工具结果处理 ContextEditingMiddleware / Tool Interceptor tool_result 返回值脱敏(PII Redaction)、长输出截断
自动摘要与压缩 SummarizationMiddleware session_before_compact, session_compact 对话历史自动压缩

概念解析:Trace与Turn

Trace(一次完整运行)

一个 Trace 是从用户按下回车、到 Agent 彻底停下来、发出 agent_end 事件的整个过程。一个 Trace 包含多个 Turn。

1
2
3
4
5
6
7
一个 Trace(一次 agent_start 到 agent_end)

├── Turn 1:调模型 → 模型返回 toolUse(要读文件)→ 执行 read 工具

├── Turn 2:带着工具结果再调模型 → 模型返回 toolUse(还要改文件)→ 执行 edit 工具

└── Turn 3:带着工具结果再调模型 → 模型返回 stop(改好了,没有工具调用)→ agent_end

Turn(一个轮次)

一个 Turn 的定义非常精确:一次模型调用 + 这次调用触发的所有工具执行。

每个 Turn 由一对 turn_startturn_end 事件包裹。关键点:一个 Turn 只有一次模型调用。 模型返回了 toolUse → 执行那批工具 → 发送 turn_end → 这个 Turn 就结束了。把工具结果喂回去再调模型,那是下一个 Turn

所以 Trace 和 Turn 的关系就是

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
Trace(一次完整运行)
│ agent_start

├── Turn 1
│ │ turn_start
│ ├── 调模型 → toolUse → 执行工具(read + grep)
│ │ turn_end
│ │
├── Turn 2
│ │ turn_start
│ ├── 调模型 → toolUse → 执行工具(edit)
│ │ turn_end
│ │
├── Turn 3
│ │ turn_start
│ ├── 调模型 → stop → 没有工具
│ │ turn_end
│ │
│ agent_end

注意:首轮 Turn 的 turn_start 是在 runAgentLoop() 入口就发出的,然后 runLoop() 内用 firstTurn 标志跳过首圈的 turn_start,避免重复。

流程全景

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
你按下回车:"帮我读一下 src/main.ts"

│ ① 你的输入变成一条消息

UserMessage { role: "user", content: "帮我读一下 src/main.ts" }

│ ② 进入循环(agentLoop 入口)—— agent_start(一个 Trace 开始了)

└── runLoop()

│ ③ 消息转换(AgentMessage → LLM 认识的 Message)

│ ┌── Turn 1 ──────────────────────────────────────────┐
│ │ turn_start │
│ │ ④ 调用 Model(每 Turn 仅一次模型调用) │
│ │ streamSimple(model, { systemPrompt, messages }) │
│ │ ↓ 逐 token 流式返回 │
│ │ AssistantMessage { │
│ │ content: [ ..., ToolCall { name: "read", ... } ],│
│ │ stopReason: "toolUse" ← 有工具调用,继续转 │
│ │ } │
│ │ ⑤ 执行 Tool(工具的五步管道,详见第5章) │
│ │ ToolResultMessage { content: [{ text: "文件内容" }] }│
│ │ turn_end │
│ └─────────────────────────────────────────────────────┘

│ 循环判断:stopReason 是 toolUse → hasMoreToolCalls = true → 继续

│ ┌── Turn 2 ──────────────────────────────────────────┐
│ │ turn_start │
│ │ ⑥ 第二次调用 Model(工具结果已追加到消息列表) │
│ │ streamSimple(model, { messages: [..., toolResult] })│
│ │ ↓ 模型看到文件内容,开始解释 │
│ │ AssistantMessage { │
│ │ content: [ TextContent { text: "这个文件..." } ],│
│ │ stopReason: "stop" ← 没有工具调用,准备停 │
│ │ } │
│ │ turn_end │
│ └─────────────────────────────────────────────────────┘

│ 循环判断:hasMoreToolCalls = false,pendingMessages 为空
│ → 内层循环退出
│ → 外层循环检查 followUp → 空 → 外层循环退出

└── agent_end(一个 Trace 结束,共 2 个 Turn)

项目结构设计

当前 my-pi-agent 严格遵循清晰解耦的三层 Python Monorepo 架构设计:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
packages/my-coding-agent          ← 【第 3 层:产品与应用层】(开箱即用代码助手)
│ • CodingAgent 门面组装
│ • 4 大工作区安全文件工具 (read / write / edit / bash)
│ • FileMutationQueue 单文件细粒度并发互斥锁
│ • 原生异步 MCP 客户端扩展 (AsyncExitStack)

packages/my-agent-core ← 【第 2 层:框架核心层】(通用 Agent 运行时微内核)
│ • 纯函数无状态微内核 run_agent_loop (loop.py)
│ • 轻量 Harness 宿主外壳 Agent (支持 prompt_stream 原生事件流)
│ • 对话转录本自愈与断头保护引擎 (tool_history.py)
│ • 模块化会话存储子系统 session/ (9种多态实体、纯追加持久化)
│ • 12 个生命周期事件与五大决策拦截点 (events & hooks)
│ • 统一 Todo 看板与 BackgroundRunner 进程树强杀
│ • 4 层 Cheap-first 上下文压缩管线 (L3➔L1➔L2➔L4)
│ • 声明式 Skills、Subagents 多智能体与 Plugin 系统

packages/my-agent-llm ← 【第 1 层:模型边界层】(底层网络与协议转换地基)
• 统一 LLM 门面 (chat / stream / achat / achat_stream)
• 三大 Provider (OpenAI / DeepSeek / Anthropic)
• 流式 Tool Calls 增量聚合与 Token Usage 锚定

项目架构

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
┌─────────────────────────────────────────────────────────────────────────┐
│ 【1. 宿主外壳 AgentHarness】 (agent.py 瘦身 260+ 行) │
│ • 纯粹的状态与队列容器,对外暴露极窄接口 │
│ • 一等公民事件流:async def prompt_stream() -> AsyncIterator[Event] │
│ • 观察者订阅 API:subscribe(listener) -> unsubscribe 句柄 │
│ • 经典门面:run() 退化为仅 10 行消费 prompt_stream 的便利包装 │
└────────────────────────────┬────────────────────────────────────────────┘
│ 驱动

┌─────────────────────────────────────────────────────────────────────────┐
│ 【2. 纯函数无状态微内核】 (loop.py: run_agent_loop) │
│ • 纯异步生成器,与类状态彻底解耦 │
│ • 内层负责 ReAct 微观步骤 + steering 动态转向 │
│ • 外层负责宏观任务流转 + follow-up 自动收割 │
│ • CancellationToken 协作式取消,成对发射标准事件流 │
└──────────────┬──────────────────────────────────────────┬───────────────┘
│ 前置清洗 │ 数据持久化
▼ ▼
┌──────────────────────────────┐ ┌────────────────────────────────────────┐
│ 【3. 转录本自愈与防400引擎】 │ │ 【4. 模块化会话存储子系统】 (session/) │
│ • _provider_context 清洗 │ │ • entries.py : 9 种多态 Pydantic v2 │
│ 剥离空中断失败轮次 │ │ • tree.py : 纯内存 DAG 算法 (防环) │
│ • tool_history.py 自愈 │ │ • memory.py : SessionState 纯函数折叠│
│ 三阶段状态机合成中断 │ │ • storage.py : 只追加纯异步协议+内存驱动│
│ 彻底消灭断头 API 400 死锁│ │ • jsonl.py : 追加存储、文件锁与自愈 │
└──────────────────────────────┘ └────────────────────────────────────────┘

为什么使用ts而不是py

1. 行业现状:为什么标杆项目(Pi、Claude Code)优先选择 TS?

在当前的 AI Agent 工业界,像 Pi(earendil-works/pi)和 Claude Code 官方首先选择 TypeScript / Node.js,核心原因在于: 1. 天然的单线程异步事件循环:V8 引擎和 Node.js 原生基于事件循环,没有 Python GIL(全局解释器锁)的历史包袱,做异步 I/O、流式打字机和并发工具调度非常顺手; 2. Web 与终端生态丰富:TypeScript 拥有成熟的前端组件生态,方便直接在同一个语言生态下构建跨平台终端(Ink/TUI)或 Web 界面; 3. 强类型元数据(TypeScript 类型系统):其结构化类型系统在编写复杂的泛型管道和事件分发时极为灵活。

2. 我们的选择:为什么要在 Python 中从零手搓出相同水准的 Agent?

尽管 TypeScript 很流行,但在真实的 AI 研发、数据分析与算法工程界,Python 依然是绝对无可撼动的“第一公民”语言

市面上大多数 Python Agent 框架(如早期的 LangChain 等)充斥着重型抽象与黑盒嵌套,一旦遇到并发、流式截断或会话分叉就束手无策。我们发起 my-pi-agent 的使命,正是为了证明:

利用现代 Python 3.11+ 的原生异步生态(asyncio + Pydantic v2 + 结构化多态联合体),完全可以从零构建出一套与 Pi / Tau 具有同等甚至更高性能、架构优雅、零过度设计的纯粹 Agent 运行时!

  • 纯原生异步(Native Asyncio):消除传统 Python 线程锁(GIL)死锁风险,通过单线程协程与操作系统内核 I/O 多路复用(IOCP / epoll),轻松调度上百个并发后台命令;
  • 微内核纯函数化:将 ReAct 循环抽离为无状态异步生成器 run_agent_loop,事件流作为一等公民,外部消费流畅如丝;
  • 工业级自愈与防僵尸:自研 tool_history.py 三阶段状态机消除断头 400 校验死锁,全平台子进程树递归强杀(taskkill /F /Tkillpg)杜绝孤儿进程;
  • 全量 100% 离线测试驱动:不依赖任何笨重第三方框架,全套 378 项离线单元测试 20 秒跑完,每个模块接缝分明、立即可用!

背景

作者指出,当时的大语言模型研究在“推理”和“行动”两个方向上各自取得了一定进展,但都存在明显的短板:

范式 代表技术 工作机制 致命缺陷 / 局限性
纯推理范式 (Reason-Only) Chain-of-Thought (CoT) 提示模型生成多步思维链来推导答案(如数学、逻辑题)。 静态黑盒,未接地(Ungrounded):仅依赖模型内部参数,无法与外部世界交互更新知识。极易导致事实幻觉(Fact Hallucination)和错误累积(Error Propagation)(见 Figure 1 (1b))。
纯行动范式 (Act-Only) 机器人规划、Web 导航 agent 将环境观测转化为文本,利用 LLM 先验生成具体动作计划。 缺乏高层规划与工作记忆:不使用语言进行抽象高层目标推理,无法维护工作记忆,难以在复杂交互中维持长

ReAct 范式的提出

为了打破这种割裂,作者提出了 ReActReason + Act):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
┌─────────────────────────┐
│ Reasoning (Thought) │
│ - 制定/调整高层计划 │
│ - 追踪进度 & 异常处理 │
└────────────┬────────────┘
│▲
Reason to ││ Act to
Act ││ Reason
▼│
┌─────────────────────────┐
│ Action & Observation │
│ - 检索外部知识库 │
│ - 与物理/网页环境交互 │
└─────────────────────────┘
  • Reason to Act(推理指导行动):利用思考轨迹(Thoughts)来创建、维护和调整高层的行动计划。
  • Act to Reason(行动支持推理):利用行动(Actions)与外部环境(如 Wikipedia、网页环境)交互,将获取的新信息反馈给后续推理步骤。
image-20260729122948649

1. 上半部分:Hotpot QA(知识密集型问答任务)

  • 任务问题:“除了 Apple Remote,还有什么设备可以控制 Apple Remote 最初设计用来交互的那个程序?”
方法范式 表现与执行路径 结果与失败原因分析
(1a) Standard (标准回答) 直接预测:iPod 错误 ×:模型无法直接凭空猜出多步跳跃(Multi-hop)的复杂事实。
(1b) CoT (纯推理/思维链) 思考:“Apple Remote 最初是用来控制 Apple TV 的,而 Apple TV 可以被 iPhone/iPad 控制……” 错误 ×(产生幻觉):模型依靠内部知识“闭门造车”,第一步就记错了事实(其实是 Front Row 程序),导致错误传导与严重幻觉
(1c) Act-Only (纯行动无思考) Act 1: 搜索 Apple Remote 观察到程序叫 Front Row; Act 2: 搜索 Front Row 提示未找到; Act 3/4: 盲目搜索词项并直接结束。 错误 ×(陷入盲目):缺乏高层推理指导。当搜索遇到阻碍(搜索不到精确词项)时,无法分析原因,也不知道如何修正搜索关键词
(1d) ReAct (推理 + 行动) Thought 1: 决定先搜 Apple Remote 找程序名; Act 1 + Obs 1: 查到程序是 Front RowThought 2 + Act 2 + Obs 2: 搜 Front Row 没找到,思考分析原因Thought 3 + Act 3: 动态修正关键词搜 Front Row (software),成功查到; Thought 4 + Act 4: 提炼出答案并输出 keyboard function keys 正确 (闭环纠错):用思考指导检索方向,用检索结果验证并矫正思考,完美避开了幻觉与盲目。

2. 下半部分:ALFWorld(交互式具身决策任务)

  • 任务目标:“把胡椒粉瓶(pepper shaker)拿到抽屉(drawer)里。”
方法范式 表现与执行路径 结果与失败原因分析
(2a) Act-Only (纯行动) Act 1-3: 去开抽屉 1,又跑去水槽 1; Act 4-5: 在没有看到胡椒瓶的情况下,强行尝试执行 Take peppershaker 1 错误 ×(无效交互):系统不断提示 Nothing happens。模型失去了工作记忆,既不知道物品常识分布在哪,也记不住自己当前做到了哪一步
(2b) ReAct (推理 + 行动) Act 1 (Think): 引入常识推理——“胡椒瓶最可能出现在柜子(1-6)或台面(1-3)”; Act 2-7: 有针对性地逐个检查,顺利在台面 3 找到并拿到胡椒瓶; Act 8 (Think): 追踪进度——“现在已拿到胡椒瓶,下一步该去放进抽屉 1”; Act 9-11: 前往 drawer 1,打开并成功放入。 正确 (高效规划):推理提供了常识先验(去哪里找)和状态跟踪(当前完成了什么子目标),大幅提升了长流程决策的成功率。

知识密集型推理任务

什么是“知识密集型推理任务”?

我们可以把这种任务拆成两部分来理解:

  1. “知识密集”(Knowledge-Intensive)
    • 意味着解决任务极其依赖海量、精准的外部事实知识(如具体的人物生卒年份、地理位置、历史事件、作品名称等)。
    • 单凭模型自身的参数记忆(内部知识),极易因记忆模糊或知识过期而产生事实幻觉(Hallucination)
  2. “推理”(Reasoning)
    • 意味着单靠简单的“关键词检索”或“直接查一次百科”是拿不到答案的。
    • 它需要模型具备多步骤的逻辑拆解、信息聚合、常识判断与比较的能力(如:先查出 A 的身份,再根据身份去查 B,最后对比 A 和 B 的出生年份)。

两大测试基准(Domains)

1. HotpotQA(多跳问答基准)

  • 任务类型:多跳复杂问答(Multi-hop Question Answering)。

  • 核心特点

    • “多跳”(Multi-hop) 指的是解答一个问题不能只查一个实体,而是需要跨越两个或更多个 Wikipedia 页面进行“跳跃式”的信息接力。

    • 典型示例

      问题:“除了 Apple Remote,还有什么设备可以控制 Apple Remote 最初设计用来交互的那个程序?”

      • 第 1 跳(查找中间实体):先搜索 Apple Remote,查出它最初设计控制的程序叫 Front Row
      • 第 2 跳(关联推理与检索):再搜索 Front Row,查找还有哪些设备可以控制 Front Row。
      • 最终合成答案:得出键盘功能键(keyboard function keys)。
  • 论文中的实验设定(Question-only Setup)

    • 传统的 HotpotQA 评估有时会直接把相关文本段落喂给模型,但本论文采用仅给定问题(Question-only)的严苛设定。
    • 模型一开始只收到一句话问题,没有任何参考文档。它必须完全依靠 ReAct 机制,自己决定何时去搜索 Wikipedia API、搜什么词、如何根据观察结果(Observation)展开下一轮思考与搜索。

2. FEVER(事实核查与验证基准)

  • 任务类型:事实提取与验证(Fact Extraction and Verification)。

  • 核心特点

    • 三分类判断:给出一个声明(Claim),要求模型结合 Wikipedia 事实,将其分类为以下三种结果之一:

      1. SUPPORTS(支持):有明确的百科事实证明该声明为真。
      2. REFUTES(反驳):有明确的百科事实证明该声明为假。
      3. NOT ENOUGH INFO(信息不足):基于现有的百科信息无法判断真伪。
    • 对细节极度敏感

      • 典型示例

        声明:“《怪奇物语》(Stranger Things)的背景设定在印第安纳州的布卢明顿(Bloomington, Indiana)。”

        • 检索与推理:搜索 Stranger Things 后发现,剧集背景设定在印第安纳州的虚构小镇霍金斯(Hawkins, Indiana),而不是布卢明顿。
        • 最终结论:判别为 REFUTES
  • 为什么适合验证 ReAct

    • FEVER 任务中 SUPPORTSREFUTES 往往只差一两个关键词或数字细节。如果模型单凭凭记忆“硬猜”(CoT),极其容易记错或产生幻觉;只有通过 ReAct 去实时查阅 Wikipedia,才能做到精准把关。

基线(Baselines)构建

在 ReAct 论文中,CoT(思维链)和 Act-Only(纯行动)作为对比基线,确实都是通过消融(Ablation)和修改 Few-shot Prompt(少样本提示词)中的示例格式来实现的

提示词方法 Prompt 中保留的组件 实现方式(对 ReAct 轨迹的处理) 作用与定位
ReAct Question + Thought + Action + Observation + Answer 完整保留:交替生成思考与行动。 本文核心方法(推理 + 行动闭环)
CoT (Reason-Only) Question + Thought + Answer 剥离 Action & Observation:删掉所有与外部 API/环境交互的步骤,只保留纯文本推导。 纯思考基线(评估仅靠内部记忆推导的表现)
Act-Only Question + Action + Observation + Answer 剥离 Thought:删掉所有 Thought 思考语句,只保留环境动作和观察反馈。 纯行动基线(评估无规划、直接盲目操作的表现)
Standard Question + Answer 剥离 Thought, Action & Observation:删掉所有中间过程,直接给答案。 标准提示基线(直出答案)

内部与外部知识的动态融合机制(Combining Internal & External Knowledge)

作者在分析中发现:

  • ReAct 擅长结合外部知识,事实接地性(Groundedness)极强,但在遭遇复杂逻辑时,格式约束可能会稍微限制其推导灵活性。
  • COT 擅长构建流畅的逻辑推理架构,但极其依赖模型内部记忆,容易发生事实幻觉(Hallucination)

启发式退避规则(Back-off Heuristics)

作者设计了两种双向退避(Back-off)策略,由模型在运行过程中根据“置信度”或“步数限制”动态决定何时切换:

1
2
3
4
5
6
7
8
                              ┌────────────────────────┐
│ 如何协同内部与外部知识? │
└───────────┬────────────┘

┌────────────────────────────┴────────────────────────────┐
▼ ▼
【策略 A: ReAct ➔ CoT-SC】 【策略 B: CoT-SC ➔ ReAct】
(以外部检索为主,失败退回内部思考) (以内部推理为主,低信心退回外部检索)

策略 A:ReAct → CoT-SC(外部检索失败时,退回内部思考)

  • 触发条件:先使用 ReAct 尝试与外部 Wikipedia API 交互。若 ReAct 在达到最大允许步数上限(HotpotQA 设为 7 步,FEVER 设为 5 步)后仍未给出答案;
  • 动作:自动放弃 ReAct 交互,退避(Back-off)切换为 CoT-SC 模式,让大模型完全依靠自身的内部记忆与逻辑来硬拆并给出最终答案。
  • 设计逻辑:如果外部搜索陷入死胡同或没查到关键词,与其直接报空失败,不如死马当活马医,信任模型自身的内部常识。

策略 B:CoT-SC → ReAct(内部推导自信度低时,退回外部检索)

  • 触发条件:先对大模型运行 CoT-SC(采样 n 条思考轨迹,默认 n = 21 并取多数票)。如果在采样结果中,得票最高的那个答案,其票数占总采样数不到一半( < n/2
  • 动作:这说明模型内部知识非常犹豫、未能形成强共识(即内部置信度低),此时自动退避切换为 ReAct 模式,亲自去 Wikipedia 搜索验证。
  • 设计逻辑:模型内部记忆有信心(投票集中)时就直接用 CoT-SC 输出;没信心(投票分散)时才去上网查,既省时又准确。

背景

单智能体架构(single-agent architectures)面临着一个内在的优化冲突:即最大化生成响应质量(Generative response quality)与减少事实性幻觉(Mitigating factual hallucinations)之间的矛盾

既要追求生成质量(语言丰富、有文采、答复详尽),又要追求事实精确(严格受限、不乱编)。这两个目标在参数层面会产生梯度干扰(Gradient interference),对事实的过度约束往往会损害语言的表达力与实用性。

related works

大模型幻觉的定义与分类

现代分类(无源/自由问答场景)

事实性幻觉(Factual Hallucinations):生成内容与现实世界的客观事实发生矛盾(例如说“牛顿提出了相对论”)。

忠实性幻觉(Faithfulness Hallucinations):生成内容不符合用户的指令要求或上文语境逻辑(例如答非所问、前后矛盾)。

MA-CF 的针对性设计: 框架中的 幻觉分析 Agent(𝒜hallu 专门把关“事实性幻觉”,而 质量分析 Agent(𝒜qual 专门审查“忠实性幻觉”(指令对齐与逻辑一致),实现了分类施策。

单智能体缓解方案

范式分类 典型技术方法 核心思路
1. 推理时干预 (Inference-Time Intervention) Prompt 工程、思维链(CoT)、上下文感知解码(CAD) 在不改变模型参数的情况下,通过约束提示词或优化解码概率引导生成。
2. 架构与知识增强 (External Knowledge Integration) 检索增强生成(RAG) 引入外部知识库,让模型动态查询实时/真实数据,打破静态参数限制。
3. 事后参数精炼 (Post-Hoc Parameter Refinement) 领域监督微调(SFT)、参数编辑(Parameter Editing)、嵌入空间修剪 重新训练或定位修改模型内部权重神经元,纠正错误记忆。

单智能体的致命缺陷:“自环问题(Self-loop Problem)”

单模型方法让同一个 Agent 同时担任生成者(Generator)、评估者(Evaluator)和纠错者(Corrector)。这会导致:

  • 认知盲区与自我确认偏误:当模型尝试自我检查(Self-reflection)时,其评估过程依然受限于最初产生错误的同一套内部知识与启发式偏见,无法实现客观诊断。
  • 优化冲突(Trade-off):事实精准度与语言流畅度在单一模型内互相拉扯,过度限制事实会损害回答的丰富度。

多智能体系统

方案范式 代表性框架 / 文献 核心机制与思路 论文指出的内在局限性
1. 迭代辩论范式 (Iterative Debate) ChatEval (Chan et al., 2023/2024) 多个 Agent 针对同一话题开展多轮交叉辩论与质询,通过多视角沟通逐步修正错误并达成共识。 共识偏误(Consensus Bias):Agent 为了强行达成一致,常互相妥协,反而“稀释”了高质量输出或强化了表面合理的错误; ❷ 职责混杂:将事实核查与质量评估掺杂在辩论中; ❸ 高延迟与高 Token 消耗
2. 协作过滤范式 (Collaborative Filtering) AgentVerse (Chen et al., 2024) 多个 Agent 组成审查阵营,通过筛选与抑制异常离群值(Outliers),减少复杂推理中的幻觉。 盲目依赖“数量堆叠”:主要靠增加 Agent 数量来干预错误,缺乏精细的职能解耦; ❷ 重降幻、轻质量:只关注消除错误,忽视了回答的语篇表达与完整性。
3. 动态网络与角色分配 (Dynamic Agent Networks) DyLAN (Liu et al., 2023) MRBalance (Zou et al., 2025) 根据任务动态挑选/组建 Agent 团队或指定特定角色(如数据库转换、因果识别),优化特定业务流程。 缺少关注分离(Separation of Concerns):绝大多数框架依然将“事实核查”与“质量评估”混在同一个决策节点中; ❷ 无法同时兼顾事实精准度与语言表达力。

现有多 Agent 系统的三大缺陷

  1. 盲目依赖“数量堆叠(Agent Multiplicity)”:许多系统仅仅依靠增加 Agent 数量来投票减少错误,缺乏精细的职能分工。
  2. 职责混杂导致“共识偏误(Consensus Bias)”:传统辩论(Debate)方法将“查证事实”与“评估表达”混在同一讨论中。Agent 们为了达成一致共识,往往会互相妥协,最终输出被“稀释”的平庸内容,甚至强化了表面合理的错误。
  3. 重“降幻”轻“表达”:现有多 Agent 论文绝大多数只关注如何降低错误率,忽略了回答的语言质量与完整性。

Methodology

核心变量

image-20260727150801684

核心目标

将“生成一个好回答”抽象为一个寻找最优解 Af 的优化问题。

给定用户查询 Q,目标是找到一个回答 A,使得联合效用函数(Joint Utility Function)U(A|Q) 最大化:

Af = arg maxAU(A|Q) = arg maxA[w1 ⋅ F(A, Q) + w2 ⋅ Q(A, Q)]

  • F(A, Q)(Factuality / 事实性):回答在多大程度上符合真实的客观事实。
  • Q(A, Q)(Quality / 质量):回答在语言丰富度、相关性、逻辑完整性上的表现。
  • w1, w2:事实性与质量之间的潜在权重(在实际决策中由合成 Agent 动态平衡)。

MA-CF framework

image-20260727152908189

阶段 1:候选生成(Phase 1: Candidate Generation)

  • 输入:用户的原始提问 Q(包含文本与上下文环境)。
  • 执行角色初稿生成 Agent(𝒜gen
  • 处理逻辑: Agent 根据输入 Q 直接生成一份未经自我修正的初始回答草稿 Acf(Q) → Ac)。
  • 设计用意: 故意不让初稿 Agent 进行过度自我审查,目的是保留一个无偏见的原始底板,将其暴露给后续的专业审查角色,避免模型在早期通过模糊表述掩盖自己的知识盲区。

阶段 2:并行评审(Phase 2: Parallelized Critique)

初稿 Ac 和原始问题 Q 会同时被分发给两个互相独立、并行运行的诊断分支:

1. 质量分析分支(Quality Analysis Branch,绿色框)

  • 执行角色响应质量分析 Agent(𝒜qual
  • 核心指标:完整度(Completeness)、逻辑性(Logic)、缺陷分析(Deficiency)。
  • 产出质量报告 Rq = 𝒜qual(Iqual, Q, Ac)
  • 职责:评估回答是否全面回答了用户问题、逻辑是否顺畅,并提取出“好的观点(Spro)”与“逻辑缺陷(Scon)”。

2. 事实性分析分支(Factuality Analysis Branch,橙色框)

  • 执行角色幻觉分析 Agent(𝒜hallu
  • 核心指标:事实准确度(Factual Accuracy)、一致性(Consistency)。
  • 产出事实报告 Rh = 𝒜hallu(Ihallu, Q, Ac)
  • 职责:像“查事实的校对员”一样,严格审查草稿中是否存在瞎编、错漏或不符事实的内容,并输出具体的错误切片与修正说明(hi, ji)。

💡 架构亮点: 这个阶段的 keypoint 是并行(Parallelized)与解耦(Decoupled)𝒜qual 专注看“文采与逻辑”,𝒜hallu 专注查“真伪”,两者互不干扰,避免了单个 Agent 既想写好文章又想扣事实细节时的“梯度干扰”与妥协折中。

阶段 3:合成精炼(Phase 3: Synthesized Refinement)

  • 输入:汇总四个关键要素——原始问题 Q + 初始草稿 Ac + 质量报告 Rq + 事实报告 Rh

  • 执行角色最终分析与合成 Agent(𝒜synth(扮演主编/元推理者角色)。

  • 处理逻辑

    Agent 根据输入的完整上下文,执行三项精炼操作:

    1. 保留:保留草稿 Ac 中被质量报告认可且未被事实报告打掉的高质量内容。
    2. 修正:根据事实报告 Rh 的批注,对错漏切片进行精准重写与定点替换。
    3. 补充:根据质量报告 Rq 指出的逻辑漏洞或未尽事宜,补全上下文。
  • 产出最终优质回答 Af = 𝒜synth(Isynth, Q, Ac, Rq, Rh)

image-20260727165550928

Experiments

dataset

1. PreciseWiki(短文本 / 精准事实问答)

  • 评估目标:评估模型在短答案场景下的显性事实准确性(Explicit factual accuracy)
  • 主要应对幻觉:外在幻觉(Extrinsic Hallucinations),即生成内容与现实世界的客观知识或可核实事实直接冲突(如错乱的日期、实体名称、概念定义)。
  • 数据集来源与规模:源自维基百科条目,并按难度进行了分层,随机采样了 N = 2000 个样本。
  • 核心价值:作为无噪声的 Ground-truth 标准,用于评估模型的“已知与未知边界”,特别是测试模型在内部知识不足时拒绝回答不可答问题(Knowledge-aware refusal)的能力。

2. LongWiki(长文本生成 / 篇章级问答)

  • 评估目标:评估模型在长文本生成中维持长序列记忆、信息整合以及长距离语义连贯性与忠实度的能力。
  • 主要应对幻觉:隐性幻觉(Implicit Hallucinations),即模型在生成段落级长文时,为了桥接逻辑断层或补全上下文而编造细节、产生前后矛盾。
  • 数据集来源与规模:要求模型围绕复杂维基百科实体生成段落级内容,随机抽取了 N = 250 个复杂实体。
  • 核心价值:突破了传统二元(对/错)短问答的局限,测试 MA-CF 在长文展开过程中是否能稳住事实密度,避免“越写越瞎编”。

3. HaluEval 2.0(跨领域综合基准)

  • 评估目标:评估框架在不同语义领域和多元任务场景下的鲁棒性与泛化能力
  • 数据构成:HaluEval 2.0 是一个大型高质量基准(包含 3.5 万个生成与人工标注样本),覆盖三大常见任务:问答(QA)、基于知识的对话(Knowledge-Grounded Dialogue)和文本摘要(Text Summarization)
  • 采样规模:作者构建了一个分层测试集,包含跨 5 个不同领域的平衡样本 N = 800
  • 核心价值:包含大量“看似合理但实际错误(Plausible but incorrect)”的微妙幻觉样本,专门用来检验模型区分似是而非的虚构内容与真实事实的边界感知能力。

指标

1. PreciseWiki(短文本精准问答)

将模型输出划分为三个互斥集合:正确(Correct, C)、幻觉(Hallucinated, H)、拒绝回答(Rejected, R)。

  • 正确率(Correct Rate)$\frac{\vert{}C\vert{}}{N}$,评估整体答对比例。

  • 幻觉率(Hallucination Rate)$\frac{\vert{}H\vert{}}{N - \vert{}R\vert{}}$,在未拒绝的回答中出现事实错误的比例(核心安全指标)。

  • 拒答率(Rejection Rate)$\frac{\vert{}R\vert{}}{N}$,评估模型在知识不足时果断“弃答”的安全边界意识。

  • F1 分数(F1 Score):正确率与回答意愿(1 − Rejection Rate)的调和平均数,防止模型为了追求高正确率而盲目弃答或为了高回答率而胡乱猜答:

    $$\text{F1} = 2 \cdot \frac{\text{Correct Rate} \cdot (1 - \text{Rejection Rate})}{\text{Correct Rate} + (1 - \text{Rejection Rate})}$$

2. LongWiki(长文本篇章问答)

长文本不能简单套用“对/错”二元分类,作者采用了原子断言抽取协议(Atomic claim extraction protocol, 设置上限 k = 32。设 Sref 为标准答案的断言集,Sout 为模型输出抽取的断言集:

  • Recall@32(召回率)$\frac{\vert{}\text{输出中得到支持的断言}\vert{}}{\vert{}S_{ref}\vert{}}$,评估模型捕获关键细节的全面性。
  • Precision(精准率)$\frac{\vert{}\text{输出中得到支持的断言}\vert{}}{\vert{}S_{out}\vert{}}$,评估输出内容中每一个断言的事实准确性。
  • F1@32:Precision 与 Recall@32 的调和平均数,综合惩罚信息遗漏与虚构断言。

3. HaluEval 2.0(跨领域综合评估)

  • Macro Factual Rate(宏观事实率):样本级的二元评估,仅当某个样本中的所有断言都完全正确时,该样本才算事实正确
  • Micro Factual Rate(微观事实率):单个样本内部得到事实支持的断言比例(捕获样本内部的事实一致性)。
  • Average Factual Rate(平均事实率):全数据集所有样本微观事实率的平均值,代表系统的整体事实密度(Factual density)

Main results

image-20260727162153476

消融实验

变体名称 对应 Fig. A.3 模板 实验设计与假设
1. 简单增强单 Agent 图 (e) Simple Answer Without Halu 不搞多 Agent,只在单个提示词里强行要求“不胡编乱造、拒绝不懂的问题”。
2. 质量与幻觉合并 图 (b) Response Merge Agent 把“查事实”和“评质量”合并到一个通用 Agent 里。
3. 质量分析细拆分 图 (a) Separate Quality Analysis 把质量分析 Agent 再拆成“评语 Agent”和“正反论点 Agent”。
4. 仅保留质量分析 图 (c) Final Only Quality Decision 砍掉幻觉核查 Agent,只看质量分析报告。
5. 仅保留幻觉分析 图 (d) Final Only Hallucination Decision 砍掉质量分析 Agent,只看幻觉核查报告。
image-20260727165501810
image-20260727164147033

Agent 功能性分析

在 Section 4.5 中,作者做了一件更精细的事:保留 4 个 Agent 的完整架构不变,但故意“剥夺/弱化”某个 Agent Prompt 里的核心功能指令,以此观察系统的性能变化。

Llama3.1-8B-instruct 为基座,作者设计了 3 组 targeted Prompt 弱化配置(具体 Prompt 见附录 Fig. A.2):

  1. 弱化质量分析(Weaken Quality Analysis)
    • 改动:简化质量 Agent 的 Prompt,删除了要求其“权衡正反方论点(Spro, Scon”和“检查逻辑一致性”的具体指令,仅保留句式简陋的通用质量评估要求。
  2. 弱化幻觉分析(Weaken Hallucination Analysis)
    • 改动:简化幻觉 Agent 的 Prompt,仅让其查找浅层事实错误,删除了“深挖推理过程中幻觉根源/错误归因”的高阶指令。
  3. 弱化决策(Weaken Decision)
    • 改动:修改合成 Agent 的 Prompt,减少其对前两个 Agent 提交的评估报告的依赖,鼓励其做出更多自主(但缺乏依据)的独立判断。
image-20260727170347009

中间报告一致性验证

为了验证中间的质量分析 Agent(𝒜qual)和幻觉分析 Agent(𝒜hallu)出具的诊断报告到底靠不靠谱?它们自己会不会产生二次幻觉?

如果中间环节的 Agent 生成的是毫无关联的废话或错误批判(False Criticism),那么最后一步的合成 Agent 就成了“垃圾进,垃圾出”(Garbage in, Garbage out)。为了验证这两个中间 Agent 的实效性与可靠性,作者设计了一套严谨的自动化中间报告相关性评估实验

数据集抽样: 从三个数据集(PreciseWiki、LongWiki、HaluEval 2.0)中各随机抽取 20% 的样本(样本量分别为 n = 400n = 50n = 160)。

自动化“裁判”模型: 引入超大规模强推理模型 DeepSeek-R1-671B 作为自动化评估员(Evaluator)。

评估维度与标尺: 裁判模型对比 问题 Q、真实标准答案(Ground Truth)、初始草稿 AcAgent 生成的报告(RqRh,测量报告的上下文相关性与一致性(Contextual Relevance),并将得分按连续标尺划分为 3 个区间:

  • 不一致(Not Consistent,Score < 0.3:报告偏离主题、分析错误或产生二次幻觉。
  • 部分一致(Partially Consistent,0.3 ≤ Score ≤ 0.7:报告捕捉到了关键信息,但粒度或严谨度与裁判标准存在微小差异。
  • 高度一致(Highly Consistent,Score > 0.7:报告分析极其精准,完全符合裁判模型的参考标准。
image-20260727180711102

样例和错误分析

成功纠错案例

image-20260727181316288

1. 长文本案例:长尾神话知识纠偏(美索不达米亚神话中的 Udug)

  • 初始草稿(Ac:片面地将 Udug 描述为一种纯粹邪恶的恶魔(暗影、毒药、刺耳的声音)。
  • 质量 Agent(𝒜qual)诊断:指出回答是片面的,相关文献中 Udug 也具有“善良与同情”的含义,信息利用不充分。
  • 幻觉 Agent(𝒜hallu)诊断:将其归类为“部分幻觉(Partial Hallucination)”,缺乏对事物二象性(Duality)的探讨。
  • 合成 Agent(𝒜synth)精炼:保留了其与混沌黑暗相关的记载,同时补充了 Udug 善良的一面以及与女神 Ereshkigal 的联系,生成了全面且准确的最终回答。

2. 短文本案例:历史事实精确更正(三明湾号航母援救拉菲号驱逐舰)

  • 问题:4月16日,USS Shamrock Bay 派遣了什么型号的战斗机协助 USS Laffey?
  • 初始草稿(Ac:错误地回答为 "F4U Corsair"(海盗式战斗机)。
  • 质量 Agent(𝒜qual)诊断:查证历史记录指出,当时派出的实际上是 4 架 FM-2 战斗机。
  • 幻觉 Agent(𝒜hallu)诊断:识别为“事实性幻觉”,澄清 F4U 确实参与了对 Laffey 的救援(陆战队 12 架 F4U 战斗轰炸机),但 Shamrock Bay 当时具体派出的战斗机是 FM-2。
  • 合成 Agent(𝒜synth)精炼:定点更正,最终仅输出准确的答案 "FM-2"

失败案例

image-20260727181357176

1. 失败案例 1:生物学知识漏诊(带状糖蚁的繁殖行为)

  • 问题:带状糖蚁的生命周期、交配模式与蚁群结构。
  • 初始草稿:错误地称带状糖蚁“与多个雄性交配(实际上一生只交配一次,即单雄交配 monandrous)”,蚁群通常是单后制。
  • 诊断过程
    • 幻觉 Agent(𝒜hallu)成功指出了蚁群结构遗漏了“多后制(polygyny,多只蚁后共存)”的可能性。
    • 质量 Agent(𝒜qual漏诊了,未能挑战草稿中“与多个雄性交配”这一错误的繁殖行为描述。
  • 最终结果:合成 Agent 成功补充了多后制蚁群的知识,但完好地保留了关于交配模式的原始幻觉

2. 失败案例 2:统计数据虚高(足球运动员 Aleksandar Đurić 的生涯数据)

  • 问题:Đurić 为新加坡国家队取得的显著成就。
  • 初始草稿:称其出场 128 次打入 50 球(严重虚高,实际数据为 53 场打入 24 球)。
  • 诊断过程
    • 质量 Agent(𝒜qual)准确指出了草稿遗漏了 Đurić 的个人荣誉(如 AFF 最佳射手、年度最佳球员)。
    • 幻觉 Agent(𝒜hallu漏诊了,它误以为 50 球/128 场的数据已经得到验证,未能识别出这一严重的数值幻觉。
  • 最终结果:合成 Agent 补充了完整的个人荣誉,但原封不动地保留了 50 球/128 场的严重统计数据幻觉

作者归纳出了 MA-CF 的核心瓶颈

  1. “木桶效应”与能力上限: 系统的抗幻觉上限,严格取决于单个诊断 Agent 的触发敏感度(Sensitivity)。合成 Agent(𝒜synth)扮演的是“整合者”而非“全知者”,如果 𝒜qual𝒜hallu 没有把错误标注出来,合成 Agent 就无法凭空发现并纠正该遗留幻觉。
  2. 错漏传递风险
    • 幻觉 Agent 漏诊 数值/事实错误被带入最终回答;
    • 质量 Agent 漏诊 逻辑漏洞或片面观点被带入最终回答。

这为未来的改进指明了方向:要提升 MA-CF 的上限,重点在于增强中间诊断 Agent 在特定领域(如精准统计数据、复杂生物学)的核查敏感度

对比实验

image-20260727182014086

论文挑选了两种最具代表性的多智能体协作范式(均采用 Llama3.1-8B 和 Qwen3-8B 作为基座):

  1. 标准智能体辩论框架(Standard Agent Debate, SAD / ChatEval):让多个 Agent 针对问题进行多轮顺序迭代辩论(Iterative Debate),直到达成共识。
  2. 动态大模型 Agent 网络(Dynamic LLM-Agent Network, DyLAN):配置了 3 个初始 Agent 和 2 层网络结构,通过分层剪枝与早期停止机制来筛选 Agent

工程启示

实际部署与计算效率(Efficiency and Deployment Considerations)

虽然 MA-CF 性能出众,但调用 4 个 Agent 无疑会增加计算开销。作者从工程落地角度提出了极具实用价值的部署策略:

  • 成本与延迟权衡(Trade-off)
    • 相比单模型:MA-CF 调用了多次模型,成本和延迟确实高于单次生成。
    • 相比多轮辩论(Debate):由于采用了并行诊断(Parallelized Critique),MA-CF 的 Token 消耗和推理延迟比多轮辩论框架降低了 50% 以上
  • 小模型集成的硬件门槛优势
    • 部署多个 8B 参数的小模型(如 Llama3.1-8B)所需的 GPU 显存和硬件基础设施开销,远低于运行单体 100B+ 或 671B(如 DeepSeek-V3)超大模型。
  • 动态条件路由策略(Conditional Routing Strategy)
    • 作者强调:没必要让每一个简单的提问都走一遍 MA-CF 管线
    • 工程建议:在入口处加入轻量级的“置信度评估器”或分类器。简单/低风险问题直接由基座单模型回答;只有遇到复杂、高风险或容易产生幻觉的问题,才路由到 MA-CF 多 Agent 管线中。
  • 模块化即插即用(Plug-and-Play Modularity)
    • 框架具备极强的可拓展性。例如在医疗、法律、金融等专业领域,可以在不修改生成 Agent 和合成 Agent 的前提下,直接将幻觉分析 Agent(𝒜hallu)替换为领域专精微调模型外挂 RAG 检索模块

现有方法的致命缺陷

当前免训练的方法主要有对比解码(Contrastive Decoding)和注意力干预(Attention Intervention)。但作者指出,它们普遍存在一个隐式的错误假设——“同质化假设”(Homogeneous Assumption)

  • 一刀切的弊端:这些方法在整个文本生成的过程中,施加的是全局统一、上下文无关(Context-agnostic)的惩罚
  • 预算浪费:由于惩罚目标不明确,它们把有限的干预资源(干预预算)浪费在了原本“低风险”的 Token 上,导致模型的生成质量(语言流利度)与幻觉抑制之间无法达到最优平衡(比如管得太死导致模型话都不会说了)。

核心发现:幻觉的“多维异质性”

为了打破上述一刀切的局限,作者对 MLLM 的长文本解码过程进行了深度的定量分析(也就是我们刚才看的 Figure 1),并首次确凿地揭示了幻觉具有多维异质性

  • 时间轴上:后期累积。随着自回归解码的深入,幻觉在生成的中后期呈现出明显的累积趋势。
  • 语义轴上:上下文共现错觉(Contextual Co-occurrence Illusion)。当模型开始胡说八道(产生幻觉)时,它对历史生成过的“语义 Token”的注意力会异常增高。

💡 深层根源(模态鸿沟):作者进一步指出了这种现象的底层数学本质——跨模态表示空间中的模态鸿沟(Modality Gap)。在模型内部,“文本-文本”的相似度天然高于“图像-文本”的对齐度。因此,随着自回归一步步往下走,视觉信息被不断稀释,大模型开始走捷径,顺着自己前面写过的文本高密度流形一路跑偏(即语义惯性),彻底把图片抛在了脑后。

image-20260711130108252

(a) 图:生成的实体在句子中的位置分布 (Temporal Position)

这张图统计了模型生成的词在整个句子(Caption)从开头(0.0)到结尾(1.0)的位置分布。

  • 蓝色线(Factual Words / 真实词):波峰极度集中在 0.0 到 0.2 之间。这说明符合事实的物体描述往往在句子前半句就早早出现了(Factual entities earlier occurrence)。
  • 橙色线(Hallucinated Words / 幻觉词):波峰显著向右偏移,大量堆积在 0.6 到 0.9 之间。这铁证了幻觉具有明显的“后期累积”特征(Hallucinations: late-stage accumulation),模型句子越写到后面,越容易胡思乱想。

(b) 图:不同层对历史语义的注意力比例 (Attention Ratio)

这张图统计了在 Transformer 的不同层数中,模型在生成下一个词时,注意力分配给“历史文本”的比例。

  • 纵轴意义:非语法词注意力比例(Non-grammar Attention Ratio)。这里的“非语法词”指的就是具体的语义词(Semantic tokens),排除了无意义的介词、标点等语法锚点。
  • 对比结果:无论在哪一层,橙色线(幻觉词)的注意力比例都死死压着蓝色线(真实词)
  • 核心结论:当模型在产生幻觉时,它对历史生成过的“语义词”分配了异常高的注意力(Hallucinations attend more to semantic than grammar tokens)。它不看图片,而是被自己前面说过的具体词汇给吸进去了,也就是摘要里说的产生“语义惯性”。

三个最核心的维度

1. 拆分了“时间轴”:从全局通罚到后期重罚 (Temporal Decoupling)

  • 过去的方法:在生成的每一个字、每一秒,对历史文本的惩罚力度完全一样 。
  • TSAI 的拆分:它把文本生成的过程按时间步(Steps)进行了切块 。在句子刚开头的语法构建期,它不怎么插手;随着自回归越往后走,幻觉风险在后期开始累积时,它的惩罚力度才像阶梯一样层层递增(步进式渐进惩罚) 。
  • 分配结果:成功在开头保护了句子的流利度,并在后期精准扼杀了幻觉的累积 。

2. 拆分了“语义轴”:从无差别剥夺到语法保护 (Semantic Decoupling)

  • 过去的方法:一旦决定惩罚历史,就把前面生成过的所有词(不论是标点还是具体名词)的注意力全部等比例扣除。
  • TSAI 的拆分:它在历史文本的 Token 阵营里拉了一道防火墙 。它把历史词拆分成了“语法锚点(介词、标点等)”“高危历史语义词(具体的物体、实体等)” 。对于语法词它网开一面(低剥夺),只把重罚精准地拍在那些容易引发共现幻觉的“历史语义词”上 。
  • 分配结果:用最小的干预代价,击碎了引发幻觉的文本惯性,同时完全不伤及大模型说话的语言骨架 。

3. 拆分了“补偿去向”:从单向视觉注入到双特征特征补偿 (Dual-Axis Compensation)

  • 过去的方法:把扣出来的注意力额度,全部粗暴地全额塞给图像。这会导致模型出现“指令遗忘(Instruction Forgetting)”——光看图不听指挥 。
  • TSAI 的拆分:它把回收上来的“注意力资金池”进行了公允的二次切分,引入了双特征特征补偿(Dual Feature Compensation) 。这些额度会按照比例,同时分配给图像 Token(加强视觉锚定)和当前的用户指令 Token(保持任务意图)
  • 分配结果:模型不仅看图看得很准,而且牢牢记着主人的任务要求,完美解决了管得太死导致模型不听话的行业痛点 。

相关工作

1. Logit 概率分布校准流派(Logit-level & Decoding Strategies)

这一流派的核心思想是在模型预测下一个词输出的概率分布(Logits)上做文章,通过人为调整概率来扼杀幻觉

OPERA :它发现模型在产生幻觉时,注意力往往会陷入局部自我循环(Over-trust 现象) 。因此它引入了一个基于注意力聚合模式的惩罚项,并搭配了一个“回溯-分配(Retrospection-allocation)”机制 。当发现不对劲时,它会触发序列回滚(Sequence Rollback),退回去重新生成 。

VCD 和 PAI :属于对比解码(Contrastive Decoding)流派 。

  • VCD 通过给图像注入视觉噪声来制造一个“业余模型(Amateur)”,然后用原模型(Expert)的 Logits 减去业余模型的 Logits,从而过滤失真的语言先验 。
  • PAI 则更激进,它直接剥离图像输入(只给文本历史),让业余模型纯靠文本惯性盲猜,再用原模型去减它,以此抵消文本惯性 。

DoLa 和 DeCo :属于层间对比(Layer Contrastive)流派。它们利用大模型深层(包含更多事实特征)和浅层(包含更多语言先验)的 Logits 进行对比或动态修正,以此放大真实表征 。

此类前作的共同痛点: 这些 Logit 校正方法往往需要多次前向传播(例如 VCD/PAI 生成每个 Token 需要跑 2 次前向)、层级追踪序列回滚(如 OPERA),这不可避免地带来了高昂的计算和时间开销,导致推理变慢 。

2. 多阶段工作流流派(Multi-stage Pipeline Interventions)

这一流派不在单次生成里做微调,而是把任务做成了流水线 。

EAZY :在输入端下功夫,通过视觉掩码(Visual Masking)过滤掉容易引发幻觉的输入源 。

HalTrapper :设计了一套“诱导-检测-抑制”(Induce-Detect-Suppress)的完整 Pipeline,主动设套诱导模型暴露出幻觉倾向再加以拔除 。

Woodpecker :属于后处理修正(Post-hoc Correction),等模型全部说完了,再调用一个外部小工具像啄木鸟一样一行行去核对和纠错 。

此类前作的共同痛点: 架构过于流水线化(Pipeline-heavy),需要多步处理或外部调用,大幅度增加了推理延迟,无法做到真正的实时响应 。

3. 注意力矩阵干预流派(Attention Interventions)

直接修改 Transformer 内部的 post-softmax 注意力矩阵 。

AttnReal :它是 TSAI 最直接的近亲 。它同样在单次前向传播(Single Forward Pass)中运行 ,其做法是识别历史输出 Token 中的“注意力下水道(Attention Sinks)”,把这些无用权重强行抠出来挪给视觉 Token 。

TSAI 框架总览

image-20260711135118270

第一步:Token 空间划分 (Token Partitioning)

在输入端,TSAI 将上下文的所有 Token 严格划分为四个互斥的语义子集:

  1. 𝒯sys (系统 Token):如大模型固定的 Prompt(“User: ”)。
  2. 𝒯img (图像 Token):由视觉编码器转化而来的图像 Patch。
  3. 𝒯ins (指令 Token):用户输入的具体任务文本(如 “Describe the image.”)。
  4. 𝒯his (历史 Token):模型在前 t − 1 步已经自己生成出的响应文本。
    • 进一步细分为:包含实际物体的语义词(Semantic Tokens)*和用于维持结构的*语法锚点(Syntactic Anchors)

划分后,在每一步(t 步)和每一层(l 层)都能得到一根初始的注意力向量 A(l, t)

第二步:干预阶段(Intervention Stages)

这是 TSAI 最核心的两大并行/协同干预机制

1️⃣ 浅层系统注意力回收 (Shallow-Stage System Recycling)

  • 动作:针对浅层网络(sys),算法发现系统词 𝒯sys 占用了大量无用的注意力(形成 Attention Sink)。
  • 操作:通过乘以一个提取因子 ρsys,强行把分配给系统词的注意力“回收”上来,形成一个可调配的注意力预算总额 Ssys
  • 去向:将回收来的这笔额度,按照固定比例补偿分配给视觉 Token (𝒯img) 和指令 Token (𝒯ins)。

2️⃣ 宽跨度双轴历史抑制 (Broad-Span Dual-Axis Historical Suppression)

针对较深层网络(hist),为了打破幻觉的异质性,实施双轴联合绞杀

  • 时间轴(Temporal Modulation):引入步进式渐进惩罚(Progressive penalty)。随着时间步(Time step)往后推移,幻觉风险越高,惩罚力度 phist 呈阶梯状越来越重,完美贴合幻觉后期累积的节奏。
  • 语义轴(Semantic Modulation):实施差异化剥夺(Differential deprivation)。对历史 Token (𝒯his) 区别对待——如果是无意义的语法词,给予低剥夺(Low deprivation);如果是高危的历史语义词,给予高剥夺(High deprivation),从而狠狠压制语义惯性。
  • 去向:通过双轴计算,从历史词中强行扣除一笔注意力总额 Shis,同样将其补偿给经过 Prefill 阶段筛选的高显著视觉 Token (𝒯img*) 以及指令 Token (𝒯ins)。

第三步:生成与最终概率输出 (Unified Output)

  • 经过上述“扣除与补充”的动态调整后,模型得到了一个全新的、经过干预的注意力向量 (l, t)
  • 最终效果对比
    • 如果没有干预(下路):模型在预测下一个词时,受语义惯性影响,词表概率(Vocabulary prob)中幻觉词(Hallucination word,如 “ball”)的概率会反超真实词,导致输出幻觉。
    • 应用 TSAI 干预后(上路):由于高危历史语义被压制、视觉和指令被放大,词表概率得到完美校准,真实词(GT word)的概率遥遥领先,模型得以准确输出符合事实的文本。

Method

1. 问题建模与 Token 空间划分 (Problem Formulation)

在标准的 MLLM 自回归解码中,假设我们在第 t 个生成步骤、第 l 层 Transformer 网络中。此时,大模型算完 Q × K 跑完 Softmax 后,会得到当前位置对过去所有位置的初始注意力分数向量 A(l, t) ∈ ℝLL 为当前总上下文序列长度)。

为了能够“数格子”动手术,TSAI 做的第一个拆分,就是将整个上下文空间 L 划分为四个互不相交的语义子集

𝒯 = 𝒯sys ∪ 𝒯img ∪ 𝒯ins ∪ 𝒯his

  • 𝒯sys(系统 Token):视觉输入前固定的系统 Prompt 词(如 USER:)。
  • 𝒯img(图像 Token):输入的 Nv 个图像 Patch 块 。
  • 𝒯ins(指令 Token):用户输入的具体文本命令(如 Describe the image.)。
  • 𝒯his(历史 Token):截至当前步 t 之前,模型自己吐出来的所有历史生成词 。

作者指出,幻觉的本质就是注意力预算在这四个子集里分配失调,模型过于死盯着历史文本 𝒯his 产生了共现幻想,却忽略了图像事实 𝒯img

2. 浅层系统注意力回收 (Shallow-Stage System Recycling)

通过前期的定量分析(Figure 3),作者注意到在网络的指定浅层(l ∈ ℒsys,如前 10 层)中,系统词 𝒯sys 占用了极其巨大的无用注意力 。

为了“白嫖”这部分被白白浪费的预算,TSAI 在浅层网络直接祭出第一刀——硬性按比例扣除i(l, t) = (1 − ρsys)Ai(l, t),  ∀i ∈ 𝒯sys

  • ρsys ∈ (0, 1) 是预设的抽取比例 。

  • 这些被强行剥夺出来的注意力分数,会被打包累加成一个可支配的浅层系统回收资金池 Ssys(l, t)

    Ssys(l, t) = ∑i ∈ 𝒯sysρsysAi(l, t)

随后,这笔额度会通过动态比例(超参数 YinsYimg),按需补偿分配给指令 Token 和视觉 Token,从而在网络刚开始时就牢牢帮模型锚定住任务意图 。

3. 宽跨度双轴历史抑制 (Broad-Span Dual-Axis Historical Suppression)

这是 TSAI 最精妙的地方。在中深层网络中(l ∈ ℒhist),为了防止模型眼睛脱离图片、只盯着历史词瞎编,算法对历史 Token 𝒯his 进行了“时间 + 语义”的双轴细粒度拆分与重罚 :

轴一:时间轴 —— 步进式渐进惩罚 (Temporal Axis)

1. 为什么不能一刀切?

如果从句子的第一个字开始,就对历史文本施加一成不变的重罚,会带来灾难性的后果 。句子开头(如主语、谓语的构建期)大模型高度依赖前文的语法和结构,如果此时粗暴地扣掉历史注意力,大模型就会“失语”,逻辑彻底崩盘(F1分数塌陷) 。

2. 数学公式实现

作者通过 Figure 1 (a) 发现,幻觉往往在长文本生成的后半段才开始滚雪球式地爆发 。因此,作者为历史抑制设计了一个随时间步(Decoding Step t)阶梯状 monotonic 递增的动态基础惩罚因子 ρhis(t)

$$\rho_{\text{his}}^{(t)} = \min\left(\rho_{\text{max}}, \Delta\rho \cdot \lfloor\frac{t}{n}\rfloor\right) \quad \text{[cite: 692]}$$

  • $\lfloor\frac{t}{n}\rfloor$(分块下取整)n 是指定的文本块大小(Chunk size) 。这意味着惩罚不是丝滑连续变化的,而是每隔 n 个字(比如每隔 5 个字)才上一个台阶 。因为作者敏锐地洞察到,幻觉在模型内部往往是以“离散的语义块(Discrete semantic chunks)”形式累积的 。
  • Δρ:每个阶梯固定的惩罚增量 。
  • ρmax:惩罚的保护上限,防止到了生成极后期时把历史注意力彻底扣成 0,导致模型彻底忘掉前文 。

💡 时间轴的效果:

句子刚开头时(t 很小),ρhis(t) = 0,TSAI 选择冷眼旁观,不打扰大模型梳理语言结构 。随着文字越吐越多,大模型开始有编造幻觉的倾向时(t 变大),惩罚力度阶梯式暴涨,拦截网越收越紧 。

轴二:语义轴 —— 语法差异化剥夺 (Semantic Axis)

1. 拦截网里的“无辜者”

即使到了生成后期,大模型去查阅历史文本时,历史文本里也包含两类完全不同的 Token :

  • 语法锚点(Syntactic Anchors):如标点符号(,.)、连词(and)、介词(inon) 。
  • 具体语义词(Semantic Tokens):如具体的物体名词(dogfridgeball) 。

如果无差别通罚,把标点和介词的注意力也扣光,模型说话就会变得颠三倒四(比如连主谓宾都连不起来) 。真正引发幻觉惯性的,是那些具体的物体语义词

2. 数学公式实现

作者引入了一个指示函数 𝕀[i ∈ 𝒢](其中 𝒢 是一个预设的、包含标点介词的静态语法词表库) 。通过它来作为标点符号的防护盾,计算出每个历史位置 i 最终承受的 deprivation 比例 ρi(t)

ρi(t) = ρhis(t) ⋅ (1 − (1 − β) ⋅ 𝕀[i ∈ 𝒢])  [cite: 697]

  • 如果当前位置 i 对应的是一个“标点/语法词”: 此时 i ∈ 𝒢,指示函数 𝕀 = 1 。 代入公式后,括号里变成了 1 − (1 − β) = β。 最终惩罚降级为:ρi(t) = β ⋅ ρhis(t) 。 (注:β ∈ [0, 1] 是一个保留系数 ,通常设得很小。这就形成了一道保护屏障,让语法词免受重罚 )。
  • 如果当前位置 i 对应的是一个“高危历史语义词(如具体的物体名)”: 此时 i ∉ 𝒢,指示函数 𝕀 = 0 。 代入公式后,括号里变成了 1 − 0 = 1。 最终惩罚直接吃满:ρi(t) = ρhis(t) —— 全额重罚,毫不留情

💰 终局:从历史打劫,全额返还视觉与指令

通过“时间 × 语义”双轴合围后,TSAI 在当前中深层网络中,成功拦截并抢救出了一大笔原本要浪费在历史语义词上的注意力总额 Shis(l, t)

Shis(l, t) = ∑i ∈ 𝒯hisρi(t)Ai(l, t)  [cite: 708]

随后,这一大笔预算会立刻通过大一统公式(Eq. 6)执行双特征特征补偿(Dual Feature Compensation),全额原地分赃 :

  1. 分给当前的用户指令 𝒯ins:把一部分额度均分给所有指令 Token,强行锚定任务意图,防止模型看图太投入而发生“指令遗忘”
  2. 分给过滤后的视觉特征 𝒯img*:利用 PGF(预充填引导过滤) 机制 ,不把额度胡乱分给全图的背景噪声(比如空白墙壁、天空) ,而是根据初始 Prefill 阶段的眼动规律 ,只精准加到那些包含核心物体的 Top-k 高显著性视觉 Token($\mathcal{T}_{\text{img}}^\*$)的格子里

做完这个大手术后,注意力向量 (l, t) 重新归一化 ,随后乘以当前层的 V 矩阵 。

4. 统一干预公式与双特征特征补偿 (Unified Formulation)

现在,算法手里拿着两笔“巨款”:浅层从系统词抠出来的 Ssys ,以及中深层从历史语义词精准打劫来的 Shis

为了把这些额度聪明地还给图像(𝒯img)和指令(𝒯ins),作者提出了最终的注意力统一加性更新公式(针对任何目标 Token j):

$$\hat{A}_j^{(l,t)} = A_j^{(l,t)} + \frac{\Delta_{\text{ins}}^{(l,t)}}{|\mathcal{T}_{\text{ins}}|}\mathbb{I}[j \in \mathcal{T}_{\text{ins}}] + \frac{S_{\text{img,sys}}^{(l,t)}}{N_v}\mathbb{I}[j \in \mathcal{T}_{\text{img}}]\cdot\mathbb{I}[l \in \mathcal{L}_{\text{sys}}] + \frac{S_{\text{img,his}}^{(l,t)}}{|\mathcal{T}_{\text{img}}^*|}\mathbb{I}[j \in \mathcal{T}_{\text{img}}^*]\cdot\mathbb{I}[l \in \mathcal{L}_{\text{hist}}]$$

在这个最终大一统公式里,作者完成了最妙的重新分配

  1. 指令特征补偿(Δins(l, t):把两笔钱中分给指令的部分合并,无缝补偿给指令 Token,彻底防止大模型只看图而“忘记任务意图”
  2. 图像预选过滤补偿(𝒯img*:在深层往图像补注意力时,作者又贴心地加了一个 PGF(预充填引导过滤) 机制 。它不把额度胡乱分给全图的背景噪声,而是利用 Prefill 阶段的先验,只精准补偿给排名前 k高显著性物体视觉 Token(𝒯img*

更新完毕后,代码会对整根 (l, t) 向量进行最后一步数值重新归一化(Re-normalization),保证分数值加起来重新等于 1,完美输送给下一步的 V 矩阵相乘 。

Experiments

Decoding Overhead (解码开销分析)

这是 TSAI 最值得骄傲的物理防线。许多号称能缓解幻觉的方法,在实际落地时因为“推理太慢”根本没法商用。作者通过 Table 2 的前向传播次数(Nforward)对比,扒下了前作们的效率底裤 。

我们来看一下论文中 Table 2 的数据分布和背后的本质含义 :

方法 (Method) 核心范式 (Paradigm) 生成每 Token 所需前向传播次数 (Nforward)
Vanilla 标准贪婪解码 (Standard) 1
DoLa / DeCo 层间对比流派 (Layer Contrastive) 1(但需要层间隐藏状态追踪或回溯修正)
AttnReal 注意力干预流派 (Attention-based) 1
Ours (TSAI) 本文方法 (Attention-based) 1(与标准 Vanilla 完全一致,零额外开销!)
VCD / PAI / CODE 传统对比解码 (Contrastive) 2(每吐一个字都必须跑 Expert 和 Amateur 双路前向)
OPERA 回溯分配流派 (Retrospection) n ( ≥ 5)(一旦触发局部循环,需要不断把序列回滚并重新计算)

🔍 为什么 TSAI 能够死死守住 Nforward = 1 的底线?

通过我们在 Method 章节里死磕的底层逻辑,你现在看这个表一定会产生降维打击的顿悟:

  • 以 VCD、PAI 或者是 CODE 为代表的流派,因为必须要在“有图事实”和“无图/噪图先验”两个平行宇宙里反复对比 ,所以它们每吐一个字,模型都必须雷打不动地跑 2 次前向传播推理延迟直接翻倍
  • OPERA 更不用说了,遇到幻觉苗头就得序列回滚(Sequence Rollback)退回去重写 ,跑一次前向的代价甚至超过 5 次 ,极其卡顿 。

TSAI 框架是在同一个 forward 函数内部,直接在算好的 Softmax 注意力矩阵上做“原地切片、原地相加”的纯数学加减法 。它没有开辟第二条平行宇宙通道(无需跑第二次前向) ,也没有任何打断、退回和重写的动作 。因此,它完美保住了 Nforward = 1 的原生速度 。

rebuttal

针对“相比前作 AttnReal 概念创新性不足”的问题

  • 审稿人(Reviewer VX7Z, 15M9)质疑:TSAI 同样是在 Softmax 之后原地做注意力重新分配,核心贡献看起来只是公式微调,有夸大创新之嫌。
  • 你的回复:直接阐明本质差异(解耦)。AttnReal 隐式地假设 Token 在时空上是同质的,而 TSAI 首次打破了这一假设,提出了4 维上下文空间注意力分配;TSAI 首次实现了对系统词(System Tokens)*中被困注意力的回收利用,且通过消退实验证明了*双特征特征补偿显著优于单向视觉注入。

针对“单次前向传播并不等同于真正低开销”的问题

  • 审稿人(Reviewer p1hu)质疑:Table 2 只数了前向传播次数,没有实际的时间运行延迟和显存开销。
  • 你的回复:直接抛出物理硬核数据:TSAI 的逐 Token 推理延迟为 46.67 ± 2.74 毫秒,相比于原生 Baseline(31.78 ± 2.20 毫秒)属于完全可接受的微弱增加;而作为交换,显存和内存占用增加极小,仅仅增加了几十个 KB,完全不构成算力瓶颈。

针对“如果固定文本生成长度为 64,结果会怎样”的问题

  • 审稿人(Reviewer VX7Z)质疑:在不同输出长度下,模型的表现和评估指标是否依然鲁棒?
  • 你的回复:补做实验证明稳健性。即使在生成的文本被强制截断为 64 个 Token 的严苛场景下,TSAI 在 LLaVA-1.5 上的 CHAIR_S / CHAIR_I 指标依然以 27.20 / 8.69 完爆原生模型的 44.00 / 13.66,同时 F1 分数(75.15 vs 73.96)依然逆势反超。

长文本后期单调递增惩罚是否会导致模型失忆、失去连贯性

代表语言长文本连贯性与质量的综合指标 F1 Score,TSAI 依然以 75.15 逆势反超了原生模型的 73.96。这用铁打的数据直接回击了审稿人——即便在长文本生成的后期,模型的长上下文连贯性和语言组织能力不仅没有发生任何“因为打压历史而崩溃失忆”的现象,反而表现得更加稳健和清醒。

基础概念

监督学习、半监督学习和无监督学习

1. 监督学习 (Supervised Learning)

  • 核心特点:数据既有特征(Features),也有标签(Labels)。即每个 X 都有对应的 Y
  • 工作原理:模型通过学习输入 X 到输出 Y 的映射关系,来预测未知数据。
  • 经典任务
    • 分类 (Classification):预测离散值(如垃圾邮件分类、多模态暴力行为检测)。
    • 回归 (Regression):预测连续值(如股票/期货价格预测)。

2. 无监督学习 (Unsupervised Learning)

  • 核心特点:数据只有特征(Features),没有标签(Labels)。只有 X,没有 Y
  • 工作原理:模型不依赖外界引导,而是依靠算法自身去挖掘数据的内在结构、相似性或潜在模式。
  • 经典任务
    • 聚类 (Clustering):把相似的数据聚集在一起(如 K-Means、用户画像分群)。
    • 降维 (Dimensionality Reduction):压缩数据特征,去除冗余(如 PCA、t-SNE)。
    • 自监督学习 (Self-Supervised Learning):现代大模型(LLM)基座预训练的核心。本质上是从无标签数据中自己构造标签(比如掩码语言模型 Masked LM,抠掉一个词让模型预测,这个词就是标签)。

3. 半监督学习 (Semi-Supervised Learning)

  • 核心特点:介于两者之间。通常是极少量的标注数据 + 大量的无标注数据
  • 为什么需要:在实际工程中,标注数据极其昂贵且耗时,而无标注数据极易获取。半监督学习旨在利用海量无标注数据包含的分布信息,来辅助提升小规模标注数据的分类性能。
  • 常见方法
    • 伪标签 (Pseudo-Labeling):先用标注数据训练一个基础模型,用它去预测无标签数据,把置信度高的预测结果当做“真标签”喂回模型重新训练。
    • 一致性正则 (Consistency Regularization):对同一个无标签输入做不同的数据增强(如加噪声),约束模型的输出保持一致。

机器学习任务

监督学习任务

  • 分类 (Classification):目标变量是离散的标签
    • 二分类:判断邮件是否为垃圾邮件、判断文本是否包含暴力行为。
    • 多分类/多标签:图像物体识别、给一篇文章贴多个领域标签。
  • 回归 (Regression):目标变量是连续的数值
    • 预测未来走势:如房价预测、股票/期货价格预测。

无监督学习任务

  • 聚类 (Clustering):无监督地将相似样本归为一类(如 K-Means、层次聚类)。
  • 降维 (Dimensionality Reduction):在高维特征空间中提炼低维特征(如 PCA、t-SNE),常用于特征降噪和数据可视化。
  • 异常检测 (Anomaly Detection):寻找与绝大多数数据显著不同的极少数样本(如信用卡盗刷、工业设备故障检测)。

常见的损失函数

一、 回归任务(Regression Loss)

用于预测连续数值(如房价预测、期货价格走势)。

1. MSE (Mean Squared Error) 均方误差 / L2 损失

  • 公式

    $$\text{MSE} = \frac{1}{n} \sum_{i=1}^{n} (y_i - \hat{y}_i)^2$$

  • 特点:对异常值(Outliers)极其敏感。因为平方项的存在,如果一个样本预测偏了,误差会被无限放大。

  • 缺点:如果数据集中有很多脏数据(噪声),模型会被异常值带偏。

2. MAE (Mean Absolute Error) 平均绝对误差 / L1 损失

  • 公式

    $$\text{MAE} = \frac{1}{n} \sum_{i=1}^{n} |y_i - \hat{y}_i|$$

  • 特点:对异常值具有很强的鲁棒性(Robust),误差的影响是线性的。

  • 缺点:在 y =  处(即误差为0的点)不可导,在训练后期(接近最优点时)可能会在最小值附近震荡,不易收敛。

3. Huber Loss (平滑的 L1 损失)

  • 核心思想:结合了 MSE 和 MAE 的优点。
  • 工作机制:当误差较小时,使用 MSE(保证梯度平滑、快速收敛);当误差较大(遇到异常值)时,自动切换为 MAE(降低敏感度,保护模型)。

二、 分类任务(Classification Loss)

用于预测离散的类别标签。

1. 二分类交叉熵损失 (Binary Cross-Entropy, BCE)

  • 公式

    $$\text{BCE} = -\frac{1}{n} \sum_{i=1}^{n} [y_i \log(\hat{y}_i) + (1 - y_i) \log(1 - \hat{y}_i)]$$

  • 适用场景:判断题(是/否),或者多标签分类(一个样本可以同时属于多个独立类别)。

  • 搭配激活函数:网络最后一层通常搭配 Sigmoid,将输出映射到 (0, 1) 区间。

2. 多分类交叉熵损失 (Categorical Cross-Entropy, CCE)

  • 公式

    $$\text{CCE} = -\sum_{c=1}^{C} y_c \log(\hat{y}_c)$$

    (其中 C 为类别总数)

  • 适用场景:单选多分类(如经典的 MNIST 手写数字识别 0~9,或是判别一段文本是否属于暴力行为)。

  • 搭配激活函数:网络最后一层必须搭配 Softmax,使所有类别的预测概率之和等于 1。

偏差(Bias)与方差(Variance)

欠拟合对应“高偏差,低方差”。模型本身偏离了真实规律,但由于简单,在不同数据集上表现都很稳定的差。

过拟合对应“低偏差,高方差”。模型对训练集拟合得极好(偏差低),但因为对数据太敏感,换一个数据集它的预测结果就会发生剧烈抖动(方差高)。

如何应对过拟合问题?

1. 数据层面 (Data Level)

  • 增加训练数据量:这是解决过拟合最直接、最根本的手段。数据量足够大、覆盖面足够广时,模型就很难“死记硬背”噪声。
  • 数据增强 (Data Augmentation):如果无法获取新数据,可以通过对现有数据进行变换(如图像旋转、裁剪、NLP中的同义词替换、多模态中引入轻微干扰)来人为制造“新数据”,提高模型的鲁棒性。

2. 模型架构与复杂度层面 (Model Level)

  • 降低模型复杂度:减少神经网络的层数(Depth)或隐藏单元数(Width),或者在使用传统机器学习(如决策树)时限制树的深度。
  • 引入正则化 (Regularization)
    • L1 正则化 (Lasso):在损失函数中加入权重绝对值之和(λ∑|w|),会使不重要的参数变为 0,从而产生稀疏解,起到特征选择的作用。
    • L2 正则化 (Ridge / 权重衰减):在损失函数中加入权重平方和($\frac{\lambda}{2} \sum w^2$),惩罚过大的权重值,让模型的参数分布更平滑,防止个别特征权重过大。
  • Dropout (深度学习特有):在每次前向传播时,随机让一定比例(如 50%)的神经元失活(输出置零)。强制网络不能依赖某几个特定的神经元组合,从而学习到更具鲁棒性的集成特征。

3. 训练策略层面 (Optimization Level)

  • 早停法 (Early Stopping):在训练过程中同时监控验证集的 Loss。当发现训练集 Loss 还在下降,但验证集 Loss 已经开始不降反升时,立即终止训练,保存验证集效果最好的那一代参数。
  • 交叉验证 (Cross-Validation):如 K-Fold 交叉验证,确保评估结果不受单一特定划分的测试集影响,让调参和模型选择更准确。
  • 集成学习 (Ensemble Learning):将多个基模型的预测结果进行组合(如 Bagging 算法、Random Forest),利用“集体智慧”抵消单个模型可能产生的过拟合风险。

正则化

什么是正则化?

定义:正则化是指在机器学习模型的目标函数(Loss Function)中引入额外约束/惩罚项的技术。

机器学习/深度学习中常见的正则化方法

在面试中,建议将方法划分为传统数学惩罚现代深度学习策略两部分来作答。

1. 传统数学惩罚项(在原 Loss 后面直接加项)

L1 正则化 (Lasso 回归)
  • 做法:在损失函数后面加上所有权重参数的绝对值之和

    $$\text{Total Loss} = \text{Original Loss} + \lambda \sum_{j=1}^{d} |w_j|$$

  • 作用与特点:会产生稀疏解(Sparse Solution)。它会无情地将很多不重要特征的权重直接削减为 0

  • 面试金句L1 正则化自带特征选择(Feature Selection)功能。” 因为在几何上,L1 的等高线是一个带尖角的方形,原 Loss 极值通常会在坐标轴(即某个 w = 0 的地方)与它相交。

L2 正则化 (Ridge 岭回归 / Weight Decay 权重衰减)
  • 做法:在损失函数后面加上所有权重参数的平方和

    $$\text{Total Loss} = \text{Original Loss} + \frac{\lambda}{2} \sum_{j=1}^{d} w_j^2$$

  • 作用与特点:会产生平滑解。它不会把权重减到 0,而是倾向于让所有的 w尽可能地接近 0 但不等于 0(惩罚大权重)。

  • 面试金句L2 正则化让模型参数分布更均匀,避免单个特征独大。” 当输入发生轻微扰动时,因为 w 都很小,输出就不会产生剧烈震荡,从而降低了方差。

梯度下降有哪些变体?

1. 批量梯度下降 (Batch Gradient Descent, BGD)

  • 工作机制:每次更新参数时,使用整个训练集的所有样本来计算梯度。
  • 优点:梯度的计算方向非常准,由于利用了全量数据,曲线下降过程非常平滑,只要学习率合适,一定能收敛到全局最优(凸问题)或局部最优(非凸问题)。
  • 缺点太慢了! 如果数据集有几百万条,每更新一次参数都要把所有数据算一遍,算力和内存/显存根本吃不消。

2. 随机梯度下降 (Stochastic Gradient Descent, SGD)

  • 工作机制:每次更新参数时,随机抽取一个样本来计算梯度并更新。
  • 优点:计算速度极快,内存占用极小。由于单个样本具有随机性,它的梯度方向总是“晃晃悠悠”的,这种随机噪声反而有助于模型跳出某些局部最优解或鞍点
  • 缺点:准确度低。即使到了山谷底部,它也不会安分停下,而是在最低点附近剧烈震荡,很难达到完美的收敛状态。

3. 小批量梯度下降 (Mini-batch Gradient Descent)

  • 工作机制:前两者的折中。每次更新参数时,使用一小批样本(一个 Batch,如 32, 64, 128, 256)来计算梯度。
  • 现代深度学习的标准做法
    • 融合了 BGD 的稳定:利用 Batch 数据的平均梯度,方向比纯 SGD 稳定得多,曲线相对平滑。
    • 融合了 SGD 的高效:不需要载入全量数据,可以完美利用 GPU 的矩阵并行计算能力。

线性与概率模型

线性回归

1. 建立模型方程

对于一个多特征的样本 X = [x1, x2, ..., xd]T,模型通过赋予每个特征不同的权重 w,再加上一个偏置 b,来计算出预测值

 = w1x1 + w2x2 + ... + wdxd + b

2. 定义损失函数:最小二乘法 (OLS)

我们如何评价这条“线”画得好不好?标准做法是计算均方误差(MSE)。我们希望所有样本的预测值与真实值之间的平方差之和最小:

$$L(W, b) = \frac{1}{2n} \sum_{i=1}^{n} (y_i - \hat{y}_i)^2$$

注:为什么要用平方?一是为了消去正负号的影响,二是平方项在数学上处处可导,非常方便计算。

3. 参数求解(怎么找到最完美的 wb

在面试中,求解方法一定要答出以下两条完全不同的路径

  • 路径 A:闭式解(解析解)—— 矩阵直接求导

    将一整套数据集写成矩阵形式,对损失函数求偏导并直接令偏导等于 0。在数学上可以一步到位直接推导出最优解公式:

    W = (XTX)−1XTY

    • 面试追问点:只有当 XTX 满秩且可逆时,才能用这个公式。如果特征之间高度相关(多重共线性),矩阵就会不可逆。此时必须引入 L2 正则化(岭回归)来强制使其可逆。
  • 路径 B:数值解 —— 梯度下降法

    当特征维度或数据量极端庞大时(比如大模型和现代深度学习场景),矩阵求逆的计算复杂度极高(O(d3))。此时我们会转而采用梯度下降法,顺着梯度的反方向一步步更新 Wb,直到模型收敛。

线性回归(Linear Regression)逻辑回归(Logistic Regression)

维度 线性回归 (Linear Regression) 逻辑回归 (Logistic Regression)
任务本质 回归(Regression)任务。 分类(Classification)任务。
输出形式 连续值。范围为 (−∞, +∞)(如预测房价、股票价格)。 离散值/概率值。范围严格限制在 (0, 1) 之间,表示属于某一类的概率。
激活函数 无(或者说是线性的 f(x) = x)。 Sigmoid 函数(将输出映射到概率区间)。
损失函数 均方误差 (MSE) / 最小二乘法。 交叉熵损失 (Cross-Entropy) / 负对数似然。

在数学和逻辑上,逻辑回归本质上是在线性回归的基础上套了一层“外壳”

逻辑回归分别是怎样处理二分类问题和多分类问题的?

直接升级为多项逻辑回归(Softmax 回归)

这是最本质、最优雅的扩展方式。当二分类的逻辑回归遇到多分类时,Sigmoid 函数会直接升级为 Softmax 函数

机制

  1. 模型不再只有一根线性输出线,而是针对 K 个类别分别拉出 K 根线,计算出 K 个类别的得分:[z1, z2, ..., zK]

  2. 使用 Softmax 函数 将这 K 个得分进行归一化,转化为一个概率分布

    $$P(y=c|X) = \frac{e^{z_c}}{\sum_{j=1}^{K} e^{z_j}}$$

  3. Softmax 的魔法在于:所有类别算出来的概率之和严格等于 1。模型最终选择概率最大的那个类别作为预测结果。

极大似然估计(MLE)

“简单来说,极大似然估计就是利用已经发生的事实(数据),去反推最有可能导致这个事实发生的一组模型参数(权重 w)。”

  • 传统概率:已知参数(比如硬币是均匀的),去预测未来的结果(抛 10 次大概有 5 次正面)。
  • 极大似然:结果已经摆在桌子上了(抛了 10 次硬币,结果 9 次正面 1 次反面),现在我们要反推参数(这枚硬币大概率被灌铅了,正面概率 w = 90% 左右最合理)。

准备工作:设定场景与符号

假设我们手头有一个二分类数据集,样本之间是独立同分布 (i.i.d.) 的。

  • 每个样本的真实标签 yi ∈ {0, 1}
  • 模型预测它为 1 的概率为 pi(在逻辑回归中 pi = σ(wTxi))。
  • 那么,模型预测它为 0 的概率自然就是 1 − pi

我们可以把这两个情况合并成一个优雅的单一样本概率公式(伯努利分布的概率质量函数):

P(yi|xi; w) = piyi(1 − pi)1 − yi

小思考:验证一下这个公式。如果真实标签 yi = 1,代入后半部分变成 (1 − pi)0 = 1,整体就剩 pi1 = pi;同理,若 yi = 0,整体就剩 1 − pi。非常完美。

数学推导三步走

现在我们要利用这个单一样本的概率,去反推最完美的权重 w

第一步:构建似然函数 L(w) —— 求总概率

由于所有样本相互独立,这 n 个样本同时发生的“总概率”,就是把所有单一样本的概率全部乘起来

$$L(w) = \prod_{i=1}^{n} P(y_i | x_i; w) = \prod_{i=1}^{n} p_i^{y_i} (1 - p_i)^{1 - y_i}$$

这个 L(w) 就是似然函数。我们的终极目标是找到一个 w,让这个连乘的总概率最大。

第二步:取对数 ln  —— 连乘变连加

直接对连乘求导会触发数学灾难(高阶乘积求导极其复杂),而且计算机会发生浮点数下溢。所以我们两边同时取自然对数 ln 

$$\ln L(w) = \ln \left( \prod_{i=1}^{n} p_i^{y_i} (1 - p_i)^{1 - y_i} \right)$$

根据对数的性质 ln (a ⋅ b) = ln a + ln b 以及 ln (ab) = bln a,我们可以把连乘的大括号拆开,变成连加

$$\ln L(w) = \sum_{i=1}^{n} \left[ y_i \ln p_i + (1 - y_i) \ln(1 - p_i) \right]$$

这就是著名的对数似然函数 (Log-Likelihood)

第三步:求偏导并令其为 0 —— 寻找极值点

为了让总概率最大,我们需要对权重 w 求偏导。这里需要用到高等数学的链式求导法则(由于 pi 内部含有 w):

$$\frac{\partial \ln L(w)}{\partial w} = 0$$

在传统的统计学中,我们解出这个方程,得到的 w 就是极大似然估计值。

交叉验证

拆解 K 折交叉验证的工作原理(以 5 折为例)

正如你所说,它的标准执行步骤非常具有仪式感,我们可以通过图解和“轮班制”来理解:

  1. 第一步(分块):把原始数据集随机打乱,并平均分成 K 个互不重叠的块(Folds)。比如我们选 K = 5,数据集就被均分为 块1、块2、块3、块4、块5。
  2. 第二步(轮流站岗):我们要进行 5 轮训练和测试
    • 第 1 轮:拿 块1 作为验证集,剩下的 块2, 3, 4, 5 合并作为训练集。训练模型,得到一个评估得分 S1
    • 第 2 轮:拿 块2 作为验证集,剩下的 块1, 3, 4, 5 作为训练集。得到得分 S2
    • ……
    • 第 5 轮:拿 块5 作为验证集,剩下的 块1, 2, 3, 4 作为训练集。得到得分 S5
  3. 第三步(大和解):把这 5 轮算出来的得分求一个平均值 Mean(S1, S2, ..., S5)这个平均分,才是我们最终认定的模型真实泛化能力。

1. 分层 K 折交叉验证 (Stratified K-Fold)

  • 对应场景样本极度不均衡。 比如你在做多模态暴力行为检测,1 万条视频里只有 100 条是暴力的(只占 1%)。
  • 做法:如果用普通 K 折,随机盲抽可能会导致某些块(Fold)里全是正常视频,一个暴力视频都没有,模型直接学不会。分层 K 折在切块时,会强迫每一个块内部的正负样本比例,都严格保持原数据集的 1:99

2. 留一法 (Leave-One-Out, LOOCV)

  • 对应场景数据量极其稀少(比如只有几十个样本)。
  • 做法:如果总共有 N 个样本,那我们就搞 N 折交叉验证。每次只把 1 个样本扣出来当验证集,剩下的 N − 1 个全部用来训练。 这样要重复跑 N 次模型。
  • 优缺点:几乎用尽了所有数据去训练,结果最精准;但如果大模型或者数据量稍大一点,算力根本承受不起,计算量爆炸。

Ridge 回归(岭回归)Lasso 回归

一、 Ridge 回归(岭回归 / L2 正则化)

Ridge 回归是在标准线性回归的均方误差(MSE)损失函数后面,加上了权重参数的平方和(称为 L2 范数惩罚项)。

1. 数学公式

$$\text{Loss}_{\text{Ridge}} = \frac{1}{2n} \sum_{i=1}^{n} (y_i - w^T x_i)^2 + \frac{\lambda}{2} \sum_{j=1}^{d} w_j^2$$

  • λ(正则化系数)用来控制惩罚的力度。λ 越大,紧箍咒越紧。

2. 工作原理与“性格”

  • 整体平滑压制L2 惩罚项对极大的权重惩罚非常严厉(因为有平方)。这会逼迫梯度下降在更新参数时,把所有的权重 w尽可能地往 0 的方向压,但绝对不会让它们真正等于 0
  • 物理意义:它让模型的参数分布变得非常均匀且平滑,避免了单个特征独大。这样当输入数据有轻微风吹草动(噪声)时,输出不会剧烈晃动,从而降低了方差。
  • 经典作用:完美解决多重共线性问题。当特征之间高度相关时,传统的线性回归矩阵求逆会崩溃,而 Ridge 回归在数学上强行保证了逆矩阵必然存在且稳定。

二、 Lasso 回归(L1 正则化)

Lasso 回归则是在 MSE 损失函数后面,加上了权重参数的绝对值之和(称为 L1 范数惩罚项)。

1. 数学公式

$$\text{Loss}_{\text{Lasso}} = \frac{1}{2n} \sum_{i=1}^{n} (y_i - w^T x_i)^2 + \lambda \sum_{j=1}^{d} \vert{}w_j\vert{}$$

2. 工作原理与“性格”

  • 无情裁剪(稀疏解)L1 惩罚项对大权重和小权重的惩罚力度是恒定的(斜率固定)。在优化过程中,它会表现得非常无情,直接把很多不重要、或者贡献小的特征的权重 w 削减到严格的 0
  • 物理意义:训练完成后,你会得到一个非常“稀疏”的权重矩阵(里面充斥着大量的 0)。
  • 经典作用自带特征选择(Feature Selection)功能。如果你的数据有 1000 个特征,Lasso 跑完可能只有 50 个特征的 w 不为 0,剩下 950 个特征直接被它无视了。这极大提升了模型在大数据场景下的可解释性和运行效率。

贝叶斯定理

$$P(\theta\vert{}X) = \frac{P(X\vert{}\theta) \cdot P(\theta)}{P(X)}$$

  • P(θ) —— 先验概率 (Prior):在没有看到新数据之前,你对这件事情发生可能性的固有认知或主观猜测。
  • P(X|θ) —— 似然概率 (Likelihood):就是我们刚刚反复聊到的“利用概率反推权重”里的那个概率。如果我的假设 θ 是对的,那么出现眼前这批数据 X 的可能性有多大?
  • P(X) —— 边缘概率/标准化常量 (Evidence):无论你的假设是什么,这批数据 X 自身发生的总概率。在很多时候它只是一个分母,用来把结果缩放到 0~1 之间。
  • P(θ|X) —— 后验概率 (Posterior)核心目标。在看到了新数据 X 之后,我们更新过后的新认知

朴素贝叶斯

回到贝叶斯公式,我们要预测一个新样本(比如一封邮件 X)属于某个类别(比如垃圾邮件 C1)的概率:

$$P(C_1 \vert{} X) = \frac{P(X \vert{} C_1) P(C_1)}{P(X)}$$

这里的特征 X 通常包含很多个维度,比如一封邮件里包含了词汇 [x1 = “发票”, x2 = “中奖”, x3 = “点击”]

在现实生活中,这些词之间明显是有关联的。但是,如果我们要去计算它们复杂的联合概率 P(“发票”, “中奖”, “点击”|C1),在数学上需要海量的数据才能统计出来,甚至会发生维度灾难。

为了打破这个僵局,朴素贝叶斯提出了一个近乎弱智、极其天真的假设 —— 特征条件独立假设(这也就是“朴素”的由来)

“它假设所有的特征之间是完全独立的、互不影响的。”

有了这个“朴素”的假设,原本极其难算的联合概率,在数学上就可以直接简单粗暴地拆解为各自概率的连乘

P(X|C1) = P(“发票”|C1) × P(“中奖”|C1) × P(“点击”|C1)

树模型与集成学习

信息增益,信息增益率

信息增益 (Information Gain) —— 初代鼻祖(ID3 标配)

要理解信息增益,必须先了解信息熵(Entropy)。信息熵是香农提出的,用来量化数据的混乱程度

📊 数学公式

对于一个数据集 D,假设里面有 K 个类别,每个类别占的比例是 pk,那么它的信息熵为:

$$\text{Entropy}(D) = - \sum_{k=1}^{K} p_k \log_2 p_k$$

  • 当数据全是一类时(纯度最高),Entropy = 0
  • 当数据各类别均匀分布时(最混乱),Entropy 达到最大值。

信息增益就是:分裂前的总熵,减去分裂后各子节点熵的加权和

$$\text{Gain}(D, A) = \text{Entropy}(D) - \sum_{v=1}^{V} \frac{\vert{}D^v\vert{}}{\vert{}D\vert{}} \text{Entropy}(D^v)$$

💡 通俗理解

信息增益代表了“得知某个特征后,系统混乱度下降了多少”。增益越大,说明这个特征分得越好,系统越快变整齐。

信息增益率 (Gain Ratio) —— 修复补丁(C4.5 标配)

为了死死卡住 ID3 偏向多取值特征的 Bug,C4.5 引入了信息增益率

📊 数学公式

信息增益率在信息增益的分子基础上,强行除以了一个分母 —— 特征自身的内在熵(Split Info)

$$\text{Gain\_ratio}(D, A) = \frac{\text{Gain}(D, A)}{\text{SplitInfo}_A(D)}$$

其中分母 $\text{SplitInfo}_A(D) = - \sum_{v=1}^{V} \frac{\vert{}D^v\vert{}}{\vert{}D\vert{}} \log_2 \frac{\vert{}D^v\vert{}}{\vert{}D\vert{}}$

💡 通俗理解

特征的取值越多、分出来的枝丫越零碎,这个特征自身的“内在熵(分母)”就会暴涨

作为分母,它就像一个无情的惩罚项。即使一个特征(比如身份证号)的信息增益很大,但因为它的取值太多导致分母极大,最终算出来的“信息增益率”也会被狠狠地拉低。从而完美抑制了过拟合。

ID3 构造决策树的五步算法流程

步骤 1:准备输入与边界检查(递归基判定)

算法传入当前的数据集 D 和剩余的特征集 A。在开始算数学公式前,先做三项“安全检查”,看是否能直接收敛为叶子节点

  • 检查 A:如果 D 中所有样本都属于同一个类别 Ck,不用分了,直接把当前节点标记为类别 Ck 的叶子节点,返回。
  • 检查 B:如果特征集 A 已经空了(特征用完了),或者 D 中所有样本在剩下特征上的取值都一模一样(无法再分),那就“少数服从多数”,把当前节点标记为 D 中样本数最多的类别的叶子节点,返回。

步骤 2:计算当前数据集的总体混乱度(总信息熵)

如果检查通过,说明需要继续分裂。首先计算当前数据集 D信息熵 H(D),作为分裂前的基准混乱度:

$$H(D) = - \sum_{k=1}^{K} p_k \log_2 p_k$$

(其中 pk 是第 k 个类别在当前数据集中的样本占比。)

步骤 3:挑选最佳分裂特征(计算信息增益)

遍历当前剩下所有特征。对于每一个特征 g

  1. 假设按特征 g 的所有可能取值,把数据集 D 划分成了多个子集 D1, D2, ..., DV

  2. 计算划分后的条件熵(即所有子集混乱度的加权平均):

    $$H(D\vert{}g) = \sum_{v=1}^{V} \frac{\vert{}D^v\vert{}}{\vert{}D\vert{}} H(D^v)$$

  3. 算出该特征的信息增益Gain(D, g) = H(D) − H(D|g)

  4. 决策判定:对比所有特征,挑出信息增益最大的那一个特征,作为当前节点的核心分裂特征 Abest

步骤 4:长出分支,切分数据集

针对挑选出的最佳特征 Abest,它有多少个可能的取值,就从当前节点向下拉出多少个对应的子分支(多叉树结构)。根据取值将数据集 D 划分到各个子节点中。

步骤 5:递归向下构建

对每一个子节点,把已经用掉的特征 Abest 从特征集中剔除,然后将子节点的数据集和缩减后的特征集重新喂回“步骤 1”,直到所有分支都触碰到叶子节点。

决策树算法是如何应对欠拟合和过拟合的

一、 决策树如何应对【过拟合】?

过拟合在决策树上的表现是:树长得太深、太茂盛,方差(Variance)极高,完美拟合训练集噪声,测试集一塌糊涂

应对过拟合,单棵决策树的核心武器是“剪枝(Pruning)”,分为预剪枝和后剪枝:

1. 预剪枝(Pre-pruning)—— 提前叫停

在树的生长过程中,只要满足某些设定的硬性阈值,就直接强行停止分裂,让其直接退化为叶子节点。

  • 限制最大深度(max_depth:这是最直观的限制,强制树高不能超过指定的层数。
  • 限制叶子节点所需最小样本数(min_samples_leaf:如果某个分支切完后,子节点里的样本数少于 5 个,就不允许再切了。
  • 限制分裂所需最小样本数(min_samples_split:如果一个节点自己包含的样本数已经很少了,直接放弃继续往下分。
  • 设置最小信息增益阈值:如果切完这刀,信息熵降低的幅度(信息增益)达不到规定标准,说明这刀不划算,不切了。
  • 优点:计算效率极高,省时省算力。
  • 缺点:非常盲目,容易带来欠拟合(因为你不知道当前的微小增益,会不会在下一次分裂时带来巨大的纯度飙升,即“视界局限”)。

2. 后剪枝(Post-pruning)—— 斩草除根

先让整棵树憋着一股劲完全长完(直到熵归 0),然后使用验证集,从下往上审视每一个非叶子节点。如果把这个节点的子树全部砍掉、直接退化为叶子节点后,验证集的准确率没有下降甚至上升了,那就果断挥刀把子树砍掉。

  • 典型算法:代价复杂度剪枝(CCP, Cost-Complexity Pruning)。它在损失函数中加入了一个关于叶子节点个数的惩罚项:

    Rα(T) = R(T) + α|Tf|

    通过调节 α,在“训练误差 R(T)”与“树的复杂程度(叶子数 |Tf|)”之间找到完美的平衡。

  • 优点:泛化能力极强,保留了真正有效的长远规则,防过拟合效果极佳。

  • 缺点:需要把树完整建好再反向遍历,算力和时间开销非常大。

二、 决策树如何应对【欠拟合】?

欠拟合在决策树上的表现是:树长得太矮、太粗糙,偏差(Bias)极高,模型在训练集和测试集上的准确率都很低

应对欠拟合的思路非常直接 —— 为模型松绑,增加模型的表达容量

  1. 放宽剪枝限制
    • 调大最大深度 max_depth
    • 调小叶子节点或分裂所需的最小样本数(min_samples_leaf / min_samples_split)。
    • 将最小信息增益阈值直接调低或设为 0,允许模型去敏锐地捕捉更加微弱的特征变化。
  2. 特征工程升级
    • 决策树如果欠拟合,很可能是当前的自变量特征根本不足以划分出正负样本。需要引入更多的交互特征(Interaction Features)、衍生特征,或者对连续特征使用更细致的离散化方案。
  3. 更换分裂指标
    • 如果你在用初代 ID3 算法,由于它无法处理连续值和缺失值,极易在复杂任务中欠拟合。此时需要无脑升级到支持连续值和二分的 C4.5CART 算法。

Boosting 算法和 Bagging 算法

维度 Bagging (自举汇聚法) Boosting (提升法)
构建方式 并行(Parallel)。各个弱学习器之间相互独立,可以同时训练。 串行(Sequential)。各个弱学习器必须串行,后一个模型依赖前一个的结果。
核心使命 降低方差(Variance)。通过平均多个过拟合的模型来消灭过拟合。 降低偏差(Bias)。通过一轮轮纠错,强行提升模型的拟合能力(消灭欠拟合)。
数据抽取 Bootstrap 抽样(有放回的随机抽样),每个模型的样本权重完全一样。 每次使用全量数据,但根据上一轮的预测错误率,动态调整样本的权重或拟合残差
弱学习器特征 倾向于使用强学习器(如长得很深的、容易过拟合的 CART 决策树)。 倾向于使用弱学习器(如只切了一刀的、极易欠拟合的“残差小树桩” Stump)。

一、 Bagging 算法的算法流程(并行架构)

Bagging(Bootstrap Aggregating)的流程核心是“独立、并行、平均分权”。

假设我们的原始数据集为 D,包含 N 个样本,我们要构建一个包含 T 个基模型的 Bagging 集成系统:

核心步骤:

  1. 并行抽样(自举汇聚)

    启动一个循环,独立重复 T 次。每一轮都对原始数据集 D 进行 Bootstrap 抽样(即有放回的随机抽样),每次抽取 N 个样本。

    • 注:因为是有放回的,某些样本会被重复抽到,某些则抽不到。最终会得到 T 个长得互不相同、但规模一样大的子数据集 {D1, D2, ..., DT}
  2. 并行独立训练

    将这 T 个子数据集同时分发出去。并行地训练 T 个基模型(各个模型之间完全闭关锁国,不知道彼此的存在)。最终得到 T 个训练好的强基模型 {f1, f2, ..., fT}

  3. 聚合投票/平均(Aggregating)

    当来了一个新样本 x 需要预测时:

    • 分类任务:让这 T 个基模型同时对 x 进行预测,统计得票数,少数服从多数(Voted),得票最多的类作为最终输出。
    • 回归任务:让这 T 个基模型输出各自的连续值,直接取算术平均值(Averaged),作为最终输出。

二、 Boosting 算法的算法流程(串行架构)

Boosting 的流程核心是“接力、串行、动态纠错”。它不搞平行宇宙,它搞的是一代代版本的迭代演进。

同样面对包含 N 个样本的数据集 D,我们要迭代 T 轮,训练出 T 个基模型进行接力:

核心步骤:

  1. 初始化“新手包”

    • 如果是调整样本权重的流派(如 AdaBoost):给原始数据集里的每个样本都赋予一个相同的初始权重(均为 $\frac{1}{N}$),此时大家是平等的。
    • 如果是拟合残差的流派(如 GBDT):先初始化一个最简单的常数预测值(比如全量标签的平均值),计算出初始的预测误差(残差)
  2. 串行接力循环(迭代 T 轮)

    进入一个严格的前后依赖循环,从 t = 1T

    • t 步训练:根据当前这轮的样本权重分布(或者当前模型留下的残差),训练第 t 个基模型 ft这个基模型被强强要求:必须拼尽全力去拟合眼前的错题/残差
    • 计算本轮话语权:计算这个基模型 ft 在训练集上的表现(如错误率)。表现越好的模型,在最终团队里的发言权重 αt 就会被分配得越大。
    • 动态更新,为下一轮铺路
      • 权重流派:把本轮 ft 做错的样本的权重调高,做对的样本权重调低,重新打包成一份“错题重灾区数据集”喂给下一轮。
      • 残差流派:用真实值减去当前总模型的预测值,算出最新的残差,作为下一轮的目标。
  3. 加权联合表决

    最终预测新样本 x 时,不是简单的民主投票,而是采取带权重的强力组合

    最终的集成模型 F(x) 是这 T 个基模型的加权求和:

    $$F(x) = \sum_{t=1}^{T} \alpha_t f_t(x)$$

    那个话语权 αt 高的“学霸模型”说了算,话语权低的“偏科模型”只起微调辅助作用。

    随机森林算法

    随机森林算法核心流程

    假设我们的原始数据集为 D,包含 N 个样本、共 M 个特征。我们准备构建一个包含 T 棵决策树的随机森林:

    步骤 1:引入“双重随机性”并行建树(核心精髓)

    启动一个并行循环,独立构建 T 棵 CART 决策树。在构建每棵树的过程中,都要注入以下两层随机性

    1. 第一重随机(样本随机)

      利用 Bootstrap 抽样(有放回的随机抽样),从原始的 N 个样本中强行抽取 N 次,形成一个用于训练当前树的子数据集 Dt

      • 注:因为是有放回的,大约会有 36.8% 的样本永远抽不到,这部分数据被称为袋外数据(OOB, Out-of-Bag),天然可以用来做不需要交叉验证的泛化评估。
    2. 第二重随机(特征随机)

      当这棵树在向下分裂节点时,算法不从全部 M 个特征中挑选最优的,而是先随机盲选一部分特征(数量通常为 $m = \sqrt{M}$log2M),然后再从这 m 个随机候选特征中,利用基尼系数(Gini)或平方误差(MSE)选出那个最好的特征进行二叉切分。

    步骤 2:树木无限生长(不剪枝)

    让这 T 棵树在各自的随机宇宙里完全长完,直到每个叶子节点都达到极高的纯度。

    • 底层哲学:因为引入了特征和样本的双重随机性,每棵树各过拟合各的(它们的过拟合噪声是高度不相关的),所以不需要单独对树进行繁琐的剪枝。

    步骤 3:多方会审,聚合预测(Aggregating)

    当森林构建完毕,来了一个全新的测试样本 x 时,所有树木同时开工:

    • 分类任务(民主投票):让 T 棵树对 x 进行分类,统计各个类别的票数,少数服从多数,得票最高的类别即为最终预测结果。
    • 回归任务(算术平均):让 T 棵树输出各自的连续值预测,直接计算这 T 个输出值的算术平均值,作为最终的预测结果。

    Adaboost

    1. 动态调整样本权重(关注错题)

    在每一轮训练中,数据集里的每个样本都有一个权重。

    • 如果某个样本被当前的弱分类器预测错误,它的权重就会在下一轮中被大幅调高,变成“重灾区错题”;
    • 如果样本预测正确,它的权重就会被调低
    • 效果:下一轮的弱分类器被迫把绝大部分精力都放在那些“前人屡屡做错的难题”上。

    2. 动态计算模型话语权(学霸权力大)

    每个弱分类器训练完后,AdaBoost 会根据它在当前数据集上的分类错误率,为它计算一个发言权重(分类器权重 α

    • 错误率越低(表现越好),这个模型的 α 就越大,在最终决议时的话语权就重;
    • 错误率接近 50%(相当于瞎猜),它的 α 就会接近 0,几乎没有话语权。

    3. 加权投票表决(精英联合)

    最终预测时,不是简单的“少数服从多数”,而是“加权投票”。把所有弱分类器的预测结果乘以它们各自的发言权重 α 进行累加,最后看正负号或者总分决定最终类别。

无监督学习:距离度量与 K-means 聚类

KNN(K-Nearest Neighbors,K近邻)

步骤 1:设定超参数 K

在开始之前,我们需要人工指定一个整数 K(比如 K = 3K = 5)。这个 K 代表我们最终要参考多少个“最近的邻居”。

  • 注:如果是二分类任务,K 通常选奇数,防止投票时出现平局。

步骤 2:计算距离(全量大扫除)

算法会遍历训练集中的每一个样本,计算新样本 x 与训练集中各个样本之间的几何距离

  • 在连续特征空间中,最常用的是欧氏距离(Euclidean Distance)

    $$d = \sqrt{\sum_{i=1}^{n} (x_i - y_i)^2}$$

  • 在某些高维或特定业务场景下,也会使用曼哈顿距离闵可夫斯基距离

步骤 3:挑选出最近的 K 个“邻居”

将计算出的所有距离进行从小到大排序,挑选出距离最近(也就是相似度最高)的 K 个训练集样本

步骤 4:投票决议,胜者为王(Aggregating)

统计这 K 个最邻居的标签类别。

  • 分类任务:采用多数表决制(Majority Voting)。这 K 个邻居里哪种类别最多,新样本就归为哪一类(比如 5 个邻居里有 4 个是“猫”,1 个是“狗”,那新样本判定为“猫”)。
  • 回归任务:如果预测的是连续值,则直接计算这 K 个邻居标签值的算术平均值作为最终输出。

K-means(K-均值)聚类

步骤 1:初始化“圈子中心”(随机选种)

在特征空间中,随机挑选 K 个点 作为初始的聚类中心(Centroids),我们把它们记为 {μ1, μ2, ..., μK}

  • 注:最原始的做法是直接从训练集里随机盲抽 K 个样本点作为中心。

步骤 2:对样本进行“划分子集”(E 步:分配身份)

遍历数据集中的每一个样本点 xi,计算它到这 K 个聚类中心的欧氏距离

$$d = \sqrt{\sum (x_{i} - \mu_k)^2}$$

通过对比,把这个样本点分配给距离它最近的那一个聚类中心所在的簇

  • 通俗理解:每个样本点都在特征空间里找离自己最近的“组织”,并加入进去。这一步结束后台面上会诞生 K 个临时的圈子。

步骤 3:重新计算“圈子中心”(M 步:中心漂移)

对于刚刚诞生的这 K 个圈子,分别计算每个圈子内部所有成员特征的算术平均值(均值)

$$\mu_k = \frac{1}{\vert{}C_k\vert{}} \sum_{x \in C_k} x$$

将算出来的这个“几何重心”,作为全新的聚类中心。此时,这 K 个中心点会发生位置的“漂移”。

步骤 4:循环迭代,直到收敛

将更新后的新中心点重新喂回 “步骤 2”,再次让所有人重新找最近的组织分配身份,接着进 “步骤 3” 重新计算中心。

如此反复循环,直到触发以下终止条件之一:

  1. 中心点不再动了:新计算出来的均值中心和上一轮的中心完全重合(或变化小于极小阈值)。
  2. 所有人的身份固化了:连续两轮迭代中,没有任何一个样本点的簇分配发生改变。
  3. 达到了最大设定的迭代步数(防止死循环)。

SVM 深度探究与贝叶斯优化

降维

解释 PCA 算法的原理和步骤

假设我们的原始数据集为 X,包含 n 个样本,每个样本有 m 个特征。即 X 是一个 n × m 的矩阵。我们希望将其降维到 k 维(k < m)。

步骤 1:特征去中心化(中心化处理)

为了消除量纲对均值的影响,必须将每个特征的均值归零。计算每个特征(每一列)的平均值 μj,然后让每个样本的该特征都减去这个均值:

Xnew = X − μ

经过这一步后,数据集的中心点完美平移到了坐标轴的原点 (0, 0)

步骤 2:计算协方差矩阵(Covariance Matrix)

协方差矩阵用来衡量特征与特征之间的相关性。计算中心化后的矩阵 Xnew 的协方差矩阵 Σ

$$\Sigma = \frac{1}{n-1} X_{\text{new}}^T X_{\text{new}}$$

得到的 Σ 是一个 m × m 的对称矩阵。矩阵对角线上的元素是各个特征自己的方差,非对角线上的元素是特征之间的协方差。

步骤 3:特征值分解,求出特征值与特征向量

对协方差矩阵 Σ 进行矩阵特征值分解(Eigendecomposition):

Σv = λv

  • 得到 m特征值 λ1, λ2, ..., λm(代表了新坐标轴方向上的方差大小)。
  • 以及对应的 m特征向量 v1, v2, ..., vm(代表了新坐标轴的空间走向,且彼此正交)。

步骤 4:挑选前 k 个最大特征值对应的特征向量

  1. 将特征值 λ 从大到小进行排序。
  2. 挑选出最大的前 k 个特征值,并取出它们对应的特征向量 [v1, v2, ..., vk]
  3. 把这 k 个列向量纵向拼接,组合成一个投影矩阵(变换矩阵) W,它的维度是 m × k

步骤 5:矩阵相乘,完成降维投影

将去中心化后的原始数据矩阵 Xnew 与投影矩阵 W 进行矩阵乘法,得到降维后的全新数据集 Y

Y = Xnew ⋅ W

由于 Xnew 的维度是 n × mW 的维度是 m × k,相乘后得到的 Y 维度正是 n × k。降维大功告成!

0%