vault backup: 2026-04-29 18:54:02
This commit is contained in:
@@ -1,14 +1,23 @@
|
||||
---
|
||||
Description: ""
|
||||
tags: [eino, ai-development, go, quickstart]
|
||||
create time: 2026-04-29 14:30
|
||||
date: "2026-03-24"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第一章:ChatModel 与 Message(Console)
|
||||
weight: 1
|
||||
---
|
||||
|
||||
## 概述
|
||||
|
||||
本章是 Eino 快速入门系列的第一章,带你理解 Eino 的 Component 抽象设计,并通过最简代码实现一次 ChatModel 调用(支持流式输出)。你将掌握 `schema.Message` 的基本用法,为后续构建完整的 ChatWithEino Agent 打下基础。
|
||||
|
||||
---
|
||||
|
||||
## Eino 框架简介
|
||||
|
||||
> [!question] 思考一下
|
||||
> 如果你要开发一个 AI 应用,需要支持 OpenAI、Claude、豆包等多个模型,你会如何设计代码架构才能让切换模型变得简单?
|
||||
|
||||
**Eino 是什么?**
|
||||
|
||||
Eino 是一个 Go 语言实现的 AI 应用开发框架(Agent Development Kit),旨在帮助开发者快速构建可扩展、可维护的 AI 应用。
|
||||
@@ -103,8 +112,35 @@ type BaseChatModel interface {
|
||||
|
||||
本章只涉及 `ChatModel`,后续章节会逐步引入 `Tool`、`Retriever` 等 Component。
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class Component {
|
||||
<<interface>>
|
||||
}
|
||||
class ChatModel {
|
||||
<<interface>>
|
||||
+Generate()
|
||||
+Stream()
|
||||
}
|
||||
class Tool {
|
||||
<<interface>>
|
||||
+Execute()
|
||||
}
|
||||
class Retriever {
|
||||
<<interface>>
|
||||
+Retrieve()
|
||||
}
|
||||
Component <|-- ChatModel
|
||||
Component <|-- Tool
|
||||
Component <|-- Retriever
|
||||
note for ChatModel "本章重点学习"
|
||||
```
|
||||
|
||||
## schema.Message:对话的基本单位
|
||||
|
||||
> [!tip] 核心概念
|
||||
> `Message` 是 Eino 对话系统的基石。理解 Message 的结构和角色语义,是掌握整个对话流程的关键。
|
||||
|
||||
`Message` 是 Eino 里对话数据的基本结构:
|
||||
|
||||
```go
|
||||
@@ -116,6 +152,16 @@ type Message struct {
|
||||
}
|
||||
```
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
A[Message] --> B[Role]
|
||||
A --> C[Content]
|
||||
A --> D[ToolCalls]
|
||||
B --> E["system / user / assistant / tool"]
|
||||
C --> F["文本内容"]
|
||||
D --> G["工具调用指令"]
|
||||
```
|
||||
|
||||
常用构造函数:
|
||||
|
||||
```go
|
||||
@@ -127,10 +173,15 @@ schema.ToolMessage("tool result", "call_id")
|
||||
|
||||
**角色语义:**
|
||||
|
||||
- `system`:系统指令,通常放在 messages 最前面
|
||||
- `user`:用户输入
|
||||
- `assistant`:模型回复
|
||||
- `tool`:工具调用结果(后续章节涉及)
|
||||
| 角色 | 用途 | 典型位置 |
|
||||
|------|------|----------|
|
||||
| `system` | 系统指令,定义模型行为 | messages 最前面 |
|
||||
| `user` | 用户输入 | 交替出现 |
|
||||
| `assistant` | 模型回复 | 交替出现 |
|
||||
| `tool` | 工具调用结果 | 工具调用后(后续章节) |
|
||||
|
||||
> [!warning] 常见错误
|
||||
> 注意 `system` 消息应该放在 messages 数组的最前面,而不是中间或末尾。这是大多数 LLM 的要求。
|
||||
|
||||
## 前置条件
|
||||
|
||||
@@ -181,42 +232,77 @@ go run ./cmd/ch01 -- "用一句话解释 Eino 的 Component 设计解决了什
|
||||
|
||||
按执行顺序:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A[开始] --> B[创建 ChatModel]
|
||||
B --> C["构造 messages 数组"]
|
||||
C --> D[调用 Stream 方法]
|
||||
D --> E{接收数据块}
|
||||
E -->|EOF| F[结束]
|
||||
E -->|有数据| G[打印 chunk.Content]
|
||||
G --> E
|
||||
```
|
||||
|
||||
1. **创建 ChatModel**:根据 `MODEL_TYPE` 环境变量选择 OpenAI 或 Ark 实现
|
||||
2. **构造输入 messages**:`SystemMessage(instruction)` + `UserMessage(query)`
|
||||
3. **调用 Stream**:所有 ChatModel 实现都必须支持 `Stream()`,返回 `StreamReader[*Message]`
|
||||
4. **打印结果**:迭代 `StreamReader` 逐帧打印 assistant 回复
|
||||
|
||||
关键代码片段(**注意:这是简化后的代码片段,不能直接运行****,完整代码请参考** [cmd/ch01/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch01/main.go)):
|
||||
关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch01/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch01/main.go)):
|
||||
|
||||
```go
|
||||
// 构造输入
|
||||
// 1. 构造输入消息数组
|
||||
// system 消息定义模型行为,user 消息包含用户问题
|
||||
messages := []*schema.Message{
|
||||
schema.SystemMessage(instruction),
|
||||
schema.UserMessage(query),
|
||||
schema.SystemMessage(instruction), // 系统指令
|
||||
schema.UserMessage(query), // 用户输入
|
||||
}
|
||||
|
||||
// 调用 Stream(所有 ChatModel 都必须实现)
|
||||
// 2. 调用 Stream 方法(所有 ChatModel 都必须实现)
|
||||
// 返回 StreamReader,用于逐块读取流式输出
|
||||
stream, err := cm.Stream(ctx, messages)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer stream.Close()
|
||||
defer stream.Close() // 确保资源释放
|
||||
|
||||
// 3. 循环接收流式数据
|
||||
for {
|
||||
chunk, err := stream.Recv()
|
||||
chunk, err := stream.Recv() // 读取下一个数据块
|
||||
if errors.Is(err, io.EOF) {
|
||||
break
|
||||
break // 流结束
|
||||
}
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
fmt.Print(chunk.Content)
|
||||
fmt.Print(chunk.Content) // 实时打印内容(不换行)
|
||||
}
|
||||
```
|
||||
|
||||
> [!note] 流式输出的优势
|
||||
> 使用 `Stream` 而不是 `Generate`,可以让用户更早看到响应,提升交互体验。这类似于 ChatGPT 的逐字显示效果。
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Component 接口**:定义可替换、可组合、可测试的能力边界
|
||||
- **Message**:对话数据的基本单位,通过角色区分语义
|
||||
- **ChatModel**:最基础的 Component,提供 `Generate` 和 `Stream` 两个核心方法
|
||||
- **实现选择**:通过环境变量或配置切换 OpenAI/Ark 等不同实现,业务代码无需改动
|
||||
| 核心概念 | 说明 | 本章要点 |
|
||||
|----------|------|----------|
|
||||
| **Component 接口** | 定义可替换、可组合、可测试的能力边界 | 通过接口实现解耦,支持多模型切换 |
|
||||
| **Message** | 对话数据的基本单位,通过角色区分语义 | 使用构造函数创建消息 |
|
||||
| **ChatModel** | 最基础的 Component | 提供 `Generate` 和 `Stream` 方法 |
|
||||
| **实现选择** | 通过环境变量或配置切换不同实现 | 业务代码无需改动 |
|
||||
|
||||
> [!success] 学习成果
|
||||
> 完成本章后,你应该能够:
|
||||
> - 理解 Eino 的 Component 设计理念
|
||||
> - 使用 ChatModel 进行单次对话调用
|
||||
> - 掌握 Message 的构造和使用方法
|
||||
> - 运行示例代码并观察流式输出
|
||||
|
||||
## 下一章预告
|
||||
|
||||
[[Eino/quick_start/chapter_02_agent_and_runner|第二章:Agent 与 Runner]] 将引入执行抽象,实现多轮对话和会话管理。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/chapter_02_agent_and_runner]]
|
||||
- [[Eino/README]]
|
||||
|
||||
@@ -1,13 +1,18 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-12"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
tags: [eino, ai-development, go, quickstart, agent, adk]
|
||||
create time: 2026-04-29 15:00
|
||||
title: 第二章:ChatModelAgent、Runner、AgentEvent(Console 多轮)
|
||||
weight: 2
|
||||
---
|
||||
|
||||
本章目标:引入 ADK 的执行抽象(Agent + Runner),并用一个 Console 程序实现多轮对话。
|
||||
## 概述
|
||||
|
||||
在第一章掌握了 `ChatModel` 组件的基础用法后,本章引入 Eino ADK 中的执行抽象——**Agent + Runner**。通过创建一个 Console 程序实现多轮对话,你将理解 Agent 接口的设计意图、事件驱动的执行模型,以及 `AsyncIterator` 如何支持流式消费。
|
||||
|
||||
---
|
||||
|
||||
<!-- @block-anchor:overview:start -->
|
||||
<!-- @block-anchor:overview:end -->
|
||||
|
||||
## 代码位置
|
||||
|
||||
@@ -40,29 +45,37 @@ you> 再用一句话总结一下
|
||||
|
||||
第一章我们学习了 **Component**(组件),它是 Eino 中可替换、可组合的能力单元:
|
||||
|
||||
- `ChatModel`:调用大语言模型
|
||||
- `Tool`:执行特定任务
|
||||
- `Retriever`:检索信息
|
||||
- `Loader`:加载数据
|
||||
| Component | 职责 | 示例 |
|
||||
|-----------|------|------|
|
||||
| `ChatModel` | 调用大语言模型 | OpenAI、Ark、Claude |
|
||||
| `Tool` | 执行特定任务 | 文件读取、代码搜索 |
|
||||
| `Retriever` | 检索信息 | 向量检索、关键词检索 |
|
||||
| `Loader` | 加载数据 | 文档解析器 |
|
||||
|
||||
> [!question] 思考一下
|
||||
> 假设你现在有一个 `ChatModel` 和一个 `Tool`,你能独立完成一个多轮对话的 AI 助手吗?如果能,你觉得会遇到哪些挑战?
|
||||
|
||||
**Component 和 Agent 的关系:**
|
||||
|
||||
- **Component 不构成完整的 AI 应用**:它只是能力单元,需要被组织、编排、执行
|
||||
- **Agent 是完整的 AI 应用**:它封装了完整的业务逻辑,可以直接运行
|
||||
- **Agent 内部使用 Component**:最核心的是 `ChatModel`(对话能力)和 `Tool`(执行能力)
|
||||
- **Component 是积木**——单个 Component 只是能力单元,需要被组织、编排、执行
|
||||
- **Agent 是整栋建筑**——它封装了完整的业务逻辑,可以直接运行
|
||||
- **Agent 内部使用 Component**——最核心的是 `ChatModel`(对话能力)和 `Tool`(执行能力)
|
||||
|
||||
**为什么需要 Agent?**
|
||||
|
||||
如果只有 Component,你需要自己:
|
||||
如果只有 Component,你需要自己管理:
|
||||
|
||||
- 管理对话历史
|
||||
- 编排调用流程(何时调用模型、何时调用工具)
|
||||
- 处理流式输出
|
||||
- 实现中断恢复
|
||||
- 对话历史的多轮累积
|
||||
- 调用流程编排(何时调模型、何时调工具)
|
||||
- 流式输出与中断处理
|
||||
- 错误恢复和状态管理
|
||||
- ...
|
||||
|
||||
**Agent 提供了什么?**
|
||||
|
||||
> [!tip] Agent 的核心价值
|
||||
> Agent = 完整运行时 + 标准事件流 + 可扩展框架。你只需要创建 Agent,然后交给 Runner 执行,不需要关心内部细节。
|
||||
|
||||
- **完整的运行时框架**:通过 `Runner` 统一管理执行过程
|
||||
- **标准的事件流输出**:`Run() -> AsyncIterator[*AgentEvent]`,支持流式、中断、恢复
|
||||
- **可扩展能力**:可以添加 tools、middleware、interrupt 等
|
||||
@@ -70,11 +83,11 @@ you> 再用一句话总结一下
|
||||
|
||||
**本章示例:**
|
||||
|
||||
`ChatModelAgent` 是最简单的 Agent,它内部只使用了 `ChatModel`,但已经具备了 Agent 的完整能力框架。后续章节会展示如何添加 `Tool` 等更多能力。
|
||||
`ChatModelAgent` 是最简单的 Agent,它内部只使用了 `ChatModel`,但已经具备了 Agent 的完整能力框架。后续章节会逐步展示如何添加 `Tool`、middleware、interrupt 等能力。
|
||||
|
||||
### Agent 接口
|
||||
|
||||
`Agent` 是 ADK 中的核心接口,定义了智能体的基本行为:
|
||||
`Agent` 是 ADK 中的核心接口,定义了智能体的基本行为。所有类型的 Agent(ChatModelAgent、WorkflowAgent、SupervisorAgent 等)都实现这个统一接口:
|
||||
|
||||
```go
|
||||
type Agent interface {
|
||||
@@ -86,22 +99,46 @@ type Agent interface {
|
||||
}
|
||||
```
|
||||
|
||||
**接口职责:**
|
||||
> [!tip] 设计精解
|
||||
> `Run()` 的返回值是 `*AsyncIterator[*AgentEvent]`——这是一个**懒加载**的流式迭代器。调用 `Run()` 时不会立即执行,只有当你开始消费事件(调用 `events.Next()`)时,Agent 才开始运行。这让你可以在启动前先配置中间件或注入依赖。
|
||||
|
||||
- `Name()` / `Description()`:标识 Agent 的名称和描述
|
||||
- `Run()`:执行 Agent 的核心方法,接收输入消息,返回事件流
|
||||
**接口职责拆解:**
|
||||
|
||||
| 方法/字段 | 职责 | 类比 |
|
||||
|-----------|------|------|
|
||||
| `Name()` | 唯一标识 Agent | 函数名 |
|
||||
| `Description()` | 描述 Agent 功能 | 函数文档 |
|
||||
| `Run()` | 执行核心逻辑 | 函数调用 |
|
||||
|
||||
**设计理念:**
|
||||
|
||||
- **统一抽象**:所有 Agent(ChatModelAgent、WorkflowAgent、SupervisorAgent 等)都实现这个接口
|
||||
- **事件驱动**:通过事件流(`AsyncIterator[*AgentEvent]`)输出执行过程,支持流式响应
|
||||
- **可扩展性**:后续加入 tools、middleware、interrupt 等能力时,接口保持不变
|
||||
```mermaid
|
||||
flowchart LR
|
||||
A["Agent 接口"] --> B["ChatModelAgent"]
|
||||
A --> C["WorkflowAgent"]
|
||||
A --> D["SupervisorAgent"]
|
||||
A --> E["..."]
|
||||
B --> F["统一 Runner 执行"]
|
||||
C --> F
|
||||
D --> F
|
||||
G["统一抽象<br/>运行时多态"] :::noteStyle
|
||||
F -.-> G
|
||||
classDef noteStyle fill:#fff3e0,stroke:#ffb74d,stroke-width:2px;
|
||||
```
|
||||
|
||||
1. **统一抽象**:所有 Agent 类型都实现同一个接口,Runner 无需关心 Agent 内部实现
|
||||
2. **事件驱动**:通过事件流输出,支持流式响应、中断恢复、状态转移
|
||||
3. **开闭原则**:新增 Agent 类型时,Runner 和消费者代码无需修改
|
||||
|
||||
### ChatModelAgent
|
||||
|
||||
`ChatModelAgent` 是 Agent 接口的一个实现,基于 ChatModel 构建:
|
||||
|
||||
```go
|
||||
// 核心参数说明:
|
||||
// - Name / Description: Agent 的身份标识
|
||||
// - Instruction: 系统指令,定义 Agent 的行为风格和目标
|
||||
// - Model: 底层的 ChatModel 组件,负责实际的模型调用
|
||||
agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
Name: "Ch02ChatModelAgent",
|
||||
Description: "A minimal ChatModelAgent with in-memory multi-turn history.",
|
||||
@@ -112,11 +149,14 @@ agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
|
||||
**ChatModel vs ChatModelAgent:本质区别**
|
||||
|
||||
> [!question] 关键辨析
|
||||
> ChatModel 和 ChatModelAgent 看起来都在"调用模型",它们的根本区别在哪里?为什么不能直接用 ChatModel 完成所有事情?
|
||||
|
||||
<table>
|
||||
<tr><td>维度</td><td>ChatModel</td><td>ChatModelAgent</td></tr>
|
||||
<tr><td><strong>定位</strong></td><td>Component(组件)</td><td>Agent(智能体)</td></tr>
|
||||
<tr><td><strong>接口</strong></td><td><pre>Generate() / Stream()</pre></td><td><pre>Run() -> AsyncIterator[*AgentEvent]</pre></td></tr>
|
||||
<tr><td><strong>输出</strong></td><td>直接返回消息内容</td><td>返回事件流(包含消息、控制动作等)</td></tr>
|
||||
<tr><td><strong>输出</strong></td><td>直接返回消息内容</td><td>返回事件流(含消息、控制动作等)</td></tr>
|
||||
<tr><td><strong>能力</strong></td><td>单纯的模型调用</td><td>可扩展 tools、middleware、interrupt 等</td></tr>
|
||||
<tr><td><strong>适用场景</strong></td><td>简单的对话场景</td><td>复杂的智能体应用</td></tr>
|
||||
</table>
|
||||
@@ -124,20 +164,23 @@ agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
**为什么需要 ChatModelAgent?**
|
||||
|
||||
1. **统一抽象**:ChatModel 只是 Component 的一种,而 Agent 是更高层的抽象,可以组合多种 Component
|
||||
2. **事件驱动**:Agent 输出事件流,支持流式响应、中断恢复、状态转移等复杂场景
|
||||
3. **可扩展性**:ChatModelAgent 可以添加 tools、middleware、interrupt 等能力,而 ChatModel 只能调用模型
|
||||
2. **事件驱动**:Agent 输出事件流,支持流式响应、中断恢复、状态转移
|
||||
3. **可扩展性**:ChatModelAgent 可以添加 tools、middleware、interrupt 等能力
|
||||
4. **编排友好**:Agent 可以被 Runner 统一管理,支持 checkpoint、恢复等运行时能力
|
||||
|
||||
> [!tip] 类比理解
|
||||
|
||||
| ChatModel | ChatModelAgent | 现实类比 |
|
||||
|-----------|----------------|----------|
|
||||
| 数据库驱动 | 业务逻辑层 | 发动机 vs 整车 |
|
||||
| 单个乐器 | 交响乐团指挥 | 砖块 vs 建筑 |
|
||||
| API 端点 | 微服务 | 积木 vs 乐高模型 |
|
||||
|
||||
**简单来说:**
|
||||
|
||||
- **ChatModel** = "负责与大语言模型通信的组件,屏蔽不同模型提供商的差异(OpenAI、Ark、Claude 等)"
|
||||
- **ChatModelAgent** = "基于模型构建的智能体,可以调用模型,但还能做更多事"
|
||||
|
||||
**类比理解:**
|
||||
|
||||
- **ChatModel** 就像"数据库驱动":负责与数据库通信,屏蔽 MySQL/PostgreSQL 的差异
|
||||
- **ChatModelAgent** 就像"业务逻辑层":基于数据库驱动构建,但还包含业务规则、事务管理等
|
||||
|
||||
**特点:**
|
||||
|
||||
- 封装了 ChatModel 的调用逻辑
|
||||
@@ -150,18 +193,19 @@ agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
|
||||
```go
|
||||
type Runner struct {
|
||||
a Agent // 要执行的 Agent
|
||||
enableStreaming bool
|
||||
store CheckPointStore // 用于中断恢复的状态存储
|
||||
a Agent // 要执行的 Agent
|
||||
enableStreaming bool // 是否启用流式输出
|
||||
store CheckPointStore // 用于中断恢复的状态存储(后续章节)
|
||||
}
|
||||
```
|
||||
|
||||
**为什么需要 Runner?**
|
||||
> [!question] 为什么需要 Runner?
|
||||
> Agent 已经有了 `Run()` 方法,为什么还要多一层 Runner?直接调用不就好了吗?
|
||||
|
||||
虽然 Agent 提供了 `Run()` 方法,但直接调用会缺少很多运行时能力:
|
||||
|
||||
1. **生命周期管理**:Runner 管理 Agent 的启动、恢复、中断等状态
|
||||
2. **Checkpoint 支持**:配合 `CheckPointStore` 实现中断恢复(后续章节涉及)
|
||||
1. **生命周期管理**:Runner 统一管理 Agent 的启动、恢复、中断等状态
|
||||
2. **Checkpoint 支持**:配合 `CheckPointStore` 实现中断恢复(第七章详解)
|
||||
3. **统一入口**:提供 `Run()` 和 `Query()` 等便捷方法
|
||||
4. **事件流封装**:将 Agent 的事件流转换为可消费的 `AsyncIterator[*AgentEvent]`
|
||||
|
||||
@@ -170,44 +214,74 @@ type Runner struct {
|
||||
```go
|
||||
runner := adk.NewRunner(ctx, adk.RunnerConfig{
|
||||
Agent: agent,
|
||||
EnableStreaming: true,
|
||||
EnableStreaming: true, // 流式模式:逐 token 消费;设为 false 则等待全部完成
|
||||
})
|
||||
|
||||
// 方式 1:传入消息列表
|
||||
// 方式 1:传入完整消息历史(支持多轮对话)
|
||||
events := runner.Run(ctx, history)
|
||||
|
||||
// 方式 2:便捷方法,传入单个查询字符串
|
||||
events := runner.Query(ctx, "你好")
|
||||
```
|
||||
|
||||
> [!tip] EnableStreaming 的影响
|
||||
>
|
||||
> | 模式 | 表现 | 适用场景 |
|
||||
> |------|------|----------|
|
||||
> | `true` | Runner 逐 token 转发事件,用户可实时看到回复 | 终端 Console、Chat UI |
|
||||
> | `false` | Runner 等待 Agent 全部执行完毕再返回结果 | API 后端、批处理任务 |
|
||||
|
||||
**Runner 的执行流程:**
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["runner.Run() / runner.Query()"] --> B["创建 AsyncIterator"]
|
||||
B --> C["开始消费事件"]
|
||||
C --> D{"下一个事件"}
|
||||
D -->|Err| E["处理错误并退出"]
|
||||
D -->|Output| F["展示给终端/客户端"]
|
||||
F --> D
|
||||
D -->|Action| G["控制动作(中断/转移/退出)"]
|
||||
G --> D
|
||||
D -->|结束| H["迭代器关闭,消费完成"]
|
||||
```
|
||||
|
||||
### AgentEvent
|
||||
|
||||
`AgentEvent` 是 Runner 返回的事件单元:
|
||||
`AgentEvent` 是 Runner 返回的事件单元,代表执行过程中的一个**离散步骤**:
|
||||
|
||||
```go
|
||||
type AgentEvent struct {
|
||||
AgentName string
|
||||
RunPath []RunStep
|
||||
AgentName string // 当前执行的是哪个 Agent
|
||||
RunPath []RunStep // 当前执行路径(支持嵌套 Agent)
|
||||
|
||||
Output *AgentOutput // 输出内容
|
||||
Action *AgentAction // 控制动作
|
||||
Err error // 执行错误
|
||||
Output *AgentOutput // 输出内容
|
||||
Action *AgentAction // 控制动作
|
||||
Err error // 执行错误
|
||||
}
|
||||
```
|
||||
|
||||
**主要字段:**
|
||||
> [!note] 事件驱动设计
|
||||
> 与传统函数调用不同,Agent 的执行不是一次性的 `return result`,而是一系列事件的有序播放。这让你的应用可以实时感知每一个执行步骤——就像看直播而不是看录播。
|
||||
|
||||
- `event.Err`:执行错误
|
||||
- `event.Output.MessageOutput`:message 或 message stream(流式)
|
||||
- `event.Action`:中断/转移/退出等控制动作(后续章节用到)
|
||||
**三大核心字段:**
|
||||
|
||||
| 字段 | 含义 | 本章用途 | 后续章节 |
|
||||
|------|------|----------|----------|
|
||||
| `event.Err` | 执行过程中发生的错误 | 错误检测与退出 | 错误处理策略 |
|
||||
| `event.Output` | Agent 的输出结果 | 展示用户回复 | 流式消费、中间结果 |
|
||||
| `event.Action` | 控制动作(中断/转移/退出等) | —— | 第七章:Interrupt & Resume |
|
||||
|
||||
---
|
||||
|
||||
### AsyncIterator:事件流的消费方式
|
||||
|
||||
`Runner.Run()` 返回的是 `*AsyncIterator[*AgentEvent]`,这是一个非阻塞的流式迭代器。
|
||||
|
||||
**为什么用 AsyncIterator 而不是直接返回结果?**
|
||||
> [!question] 为什么用 AsyncIterator?
|
||||
> 为什么不直接返回 `[]*AgentEvent` 或者单个结果?
|
||||
|
||||
因为 Agent 的执行是**流式**的:模型逐 token 生成回复,Tool 调用穿插其中。如果等全部完成再返回,用户需要等待更长时间。`AsyncIterator` 让你可以实时消费每一个事件。
|
||||
因为 Agent 的执行是**流式**的:模型逐 token 生成回复,Tool 调用穿插其中。如果等全部完成再返回,用户需要等待更长时间。`AsyncIterator` 让你可以**实时消费**每一个事件。
|
||||
|
||||
**消费方式:**
|
||||
|
||||
@@ -220,90 +294,128 @@ for {
|
||||
if !ok {
|
||||
break // 迭代器关闭,全部事件已消费
|
||||
}
|
||||
|
||||
// 三种处理方式互斥,根据具体场景判断
|
||||
if event.Err != nil {
|
||||
// 处理错误
|
||||
// 1. 错误分支:执行出错,记录日志并决定是否继续
|
||||
log.Printf("agent error: %v", event.Err)
|
||||
break
|
||||
}
|
||||
|
||||
if event.Output != nil && event.Output.MessageOutput != nil {
|
||||
// 处理消息输出(可能是流式)
|
||||
// 2. 输出分支:收到消息内容(可能是流式分片)
|
||||
msg := event.Output.MessageOutput.Message
|
||||
fmt.Print(msg.Content)
|
||||
}
|
||||
|
||||
// 3. Action 分支:当前章用不到,后续章节(Interrupt/Resume)会深入
|
||||
// if event.Action != nil { ... }
|
||||
}
|
||||
```
|
||||
|
||||
**注意:**每次 `runner.Run()` 创建新的迭代器,消费一次后不可重复使用。
|
||||
> [!warning] 重要注意事项
|
||||
> - **每次 `runner.Run()` 创建新的迭代器**,消费一次后不可重复使用
|
||||
> - **不要忽略 `event.Err`**——Agent 内部可能静默失败(如工具执行超时)
|
||||
> - **注意 goroutine 安全**——多个消费者同时读取同一个 AsyncIterator 是不安全的
|
||||
|
||||
## 多轮对话的实现
|
||||
|
||||
本章实现的是简单的多轮对话:用户输入 → 模型回复 → 用户继续输入 → ...
|
||||
|
||||
**实现方式:**
|
||||
**核心思想:**
|
||||
|
||||
没有 tools 时,`ChatModelAgent` 在一次 `Run()` 里只会完成一轮模型调用。多轮对话是通过调用侧维护 history 实现的:
|
||||
没有 tools 时,`ChatModelAgent` 在一次 `Run()` 里只会完成一轮模型调用。多轮对话是通过**调用侧维护 history** 实现的——每次调用都把完整的对话历史传进去,让模型知道之前聊了什么。
|
||||
|
||||
1. 用 `history []*schema.Message` 保存累计对话
|
||||
2. 每次用户输入:把 `UserMessage` 追加到 history
|
||||
3. 调用 `runner.Run(ctx, history)` 得到事件流,消费得到 assistant 文本
|
||||
4. 把本轮 assistant 文本追加回 history,进入下一轮
|
||||
```mermaid
|
||||
flowchart TD
|
||||
S["初始化 history = []"] --> L["进入循环"]
|
||||
L --> U["用户输入 UserMessage"]
|
||||
U --> H1["追加到 history"]
|
||||
H1 --> R["runner.Run(ctx, history)"]
|
||||
R --> E["消费事件流"]
|
||||
E --> C{"有 Output?"}
|
||||
C -->|是| A1["收集 assistant 文本"]
|
||||
A1 --> H2["追加 AssistantMessage 到 history"]
|
||||
H2 --> L
|
||||
C -->|否/结束| OUT["退出循环"]
|
||||
|
||||
**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch02/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch02/main.go)):
|
||||
style S fill:#e1f5fe
|
||||
style OUT fill:#ffebee
|
||||
```
|
||||
|
||||
**逐步拆解:**
|
||||
|
||||
1. **用 `history []*schema.Message` 保存累计对话**——所有已发生过的消息都存这里
|
||||
2. **每次用户输入**:把 `UserMessage` 追加到 history
|
||||
3. **调用 `runner.Run(ctx, history)`**:得到完整事件流,消费得到 assistant 回复
|
||||
4. **把本轮 assistant 文本追加回 history**:进入下一轮时,模型能看到全部对话历史
|
||||
|
||||
**关键代码片段(注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch02/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch02/main.go)):
|
||||
|
||||
```go
|
||||
// history 维护完整的对话历史,容量预设 16 条消息
|
||||
history := make([]*schema.Message, 0, 16)
|
||||
|
||||
for {
|
||||
// 1. 读取用户输入
|
||||
// 1. 读取用户输入,空行表示退出
|
||||
line := readUserInput()
|
||||
if line == "" {
|
||||
break
|
||||
}
|
||||
|
||||
// 2. 追加用户消息到 history
|
||||
// 2. 将用户消息追加到 history
|
||||
// 这样模型在下一轮能"记住"之前的对话
|
||||
history = append(history, schema.UserMessage(line))
|
||||
|
||||
// 3. 调用 Runner 执行 Agent
|
||||
// 返回的事件流包含所有输出步骤(消息、工具调用等)
|
||||
events := runner.Run(ctx, history)
|
||||
|
||||
// 4. 消费事件流,收集 assistant 回复
|
||||
// 4. 消费事件流,收集 assistant 的回复内容
|
||||
content := collectAssistantFromEvents(events)
|
||||
fmt.Println("[assistant]", content)
|
||||
|
||||
// 5. 追加 assistant 消息到 history
|
||||
// 5. 将 assistant 回复也追加到 history
|
||||
// nil 表示本轮没有工具调用(后续章节会用到)
|
||||
history = append(history, schema.AssistantMessage(content, nil))
|
||||
}
|
||||
```
|
||||
|
||||
**流程图:**
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 初始化 history = [] │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 用户输入 UserMessage │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 追加到 history │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ runner.Run(history) │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 消费事件流 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 追加 AssistantMessage│
|
||||
└──────────────────────┘
|
||||
↓
|
||||
(循环继续)
|
||||
```
|
||||
> [!note] 关于 history 的内存管理
|
||||
>
|
||||
> 当前实现将所有消息保留在内存中。在实际应用中,你可能需要:
|
||||
> - 设置最大消息数量限制(如上面的 `16`)
|
||||
> - 使用摘要压缩(Summarization Middleware,第五章介绍)
|
||||
> - 使用外部存储(Memory 组件,第三章介绍)
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Agent 接口**:定义智能体的基本行为,核心是 `Run() -> AsyncIterator[*AgentEvent]`
|
||||
- **ChatModelAgent**:基于 ChatModel 实现的 Agent,提供统一的执行抽象
|
||||
- **Runner**:Agent 的执行入口,管理生命周期、checkpoint、事件流等运行时能力
|
||||
- **AgentEvent**:事件驱动的输出单元,支持流式响应和控制动作
|
||||
- **多轮对话**:通过调用侧维护 history 实现,每次 `Run()` 完成一轮对话
|
||||
| 核心概念 | 说明 | 关键要点 |
|
||||
|----------|------|----------|
|
||||
| **Agent 接口** | 定义智能体的基本行为,`Run() -> AsyncIterator[*AgentEvent]` | 统一抽象,所有 Agent 类型共享同一接口 |
|
||||
| **ChatModelAgent** | 基于 ChatModel 实现的 Agent | 最简 Agent,是后续扩展的基础 |
|
||||
| **Runner** | Agent 的执行入口 | 管理生命周期、Checkpoint、事件流封装 |
|
||||
| **AgentEvent** | 事件驱动的输出单元 | 包含 Output(消息)、Action(控制)、Err(错误) |
|
||||
| **AsyncIterator** | 流式迭代器,逐事件消费 | 实时响应,不阻塞等待全部完成 |
|
||||
| **多轮对话** | 调用侧维护 history 实现 | 每次 `Run()` 传完整历史,每轮追加新消息 |
|
||||
|
||||
> [!success] 学习成果
|
||||
> 完成本章后,你应该能够:
|
||||
> - 理解 Component 和 Agent 的本质区别
|
||||
> - 使用 `adk.NewChatModelAgent` 创建自己的 Agent
|
||||
> - 通过 Runner 执行 Agent 并消费事件流
|
||||
> - 实现基于 history 的多轮对话
|
||||
>
|
||||
> > [!tip] 动手练习
|
||||
> > 试着修改 `Instruction` 参数,给你的 Agent 设定一个角色(如"你是一个编程导师"),观察不同指令对回复的影响。这就是 Prompt Engineering 的雏形!
|
||||
|
||||
## 下一章预告
|
||||
|
||||
[[Eino/quick_start/chapter_03_memory_and_session|第三章:Memory 与 Session]] 将引入持久化存储机制,让对话历史跨进程保留,不再因为程序重启而丢失记忆。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/chapter_01_chatmodel_and_message]]
|
||||
- [[Eino/quick_start/chapter_03_memory_and_session]]
|
||||
- [[Eino/core_modules/eino_adk/agent_interface]]
|
||||
- [[Eino/core_modules/eino_adk/agent_implementation/chat_model]]
|
||||
|
||||
@@ -1,30 +1,27 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-12"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
tags: [eino, ai-development, go, quickstart, memory, session]
|
||||
create time: 2026-04-29 15:30
|
||||
title: 第三章:Memory 与 Session(持久化对话)
|
||||
weight: 3
|
||||
---
|
||||
|
||||
本章目标:实现对话历史的持久化存储,支持跨进程恢复会话。
|
||||
## 概述
|
||||
|
||||
> **⚠️ 重要说明:业务层概念 vs 框架概念**
|
||||
在第二章掌握了多轮对话的实现后,我们面临一个关键问题:**对话历史只存在于内存中,进程退出后一切归零**。本章引入 **Memory 与 Session** 机制,让对话历史能够持久化保存并跨进程恢复,为构建真正的智能助手奠定基础。
|
||||
|
||||
> 本章介绍的 **Memory、Session、Store 是业务层概念**,**不是 Eino 框架的核心组件**。
|
||||
|
||||
>
|
||||
|
||||
> 换句话说,Eino 框架只负责"如何处理消息",而"如何存储消息"完全由业务层决定。本章提供的实现只是一个简单的参考示例,你可以根据自己的业务需求选择完全不同的存储方案(数据库、Redis、云存储等)。
|
||||
> [!warning] 重要概念区分:业务层 vs 框架层
|
||||
> 本章介绍的 **Memory、Session、Store 是业务层概念**,**不是 Eino 框架的核心组件**。Eino 框架只负责"如何处理消息",而"如何存储消息"完全由业务层决定。本章提供的实现只是一个参考示例,你可以根据自己的需求选择数据库、Redis、云存储等方案。
|
||||
|
||||
## 代码位置
|
||||
|
||||
- 入口代码:[cmd/ch03/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch03/main.go)
|
||||
- Memory 实现:[mem/store.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/mem/store.go)
|
||||
- Store 实现:[mem/store.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/mem/store.go)
|
||||
|
||||
---
|
||||
|
||||
## 前置条件
|
||||
|
||||
与第一章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)。
|
||||
与第二章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)。
|
||||
|
||||
## 运行
|
||||
|
||||
@@ -53,9 +50,15 @@ Session saved: 083d16da-6b13-4fe6-afb0-c45d8f490ce1
|
||||
Resume with: go run ./cmd/ch03 --session 083d16da-6b13-4fe6-afb0-c45d8f490ce1
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
<!-- @block-anchor:overview:start -->
|
||||
<!-- @block-anchor:overview:end -->
|
||||
|
||||
## 从内存到持久化:为什么需要 Memory
|
||||
|
||||
第二章我们实现了多轮对话,但有一个问题:**对话历史只存在于内存中**。
|
||||
> [!question] 思考一下
|
||||
> 第二章我们实现了多轮对话,但存在一个问题——如果进程退出、机器重启,之前聊的内容还会在吗?
|
||||
|
||||
**内存存储的局限:**
|
||||
|
||||
@@ -71,12 +74,17 @@ Resume with: go run ./cmd/ch03 --session 083d16da-6b13-4fe6-afb0-c45d8f490ce1
|
||||
|
||||
**简单类比:**
|
||||
|
||||
- **内存存储** = "草稿纸"(进程退出就没了)
|
||||
- **Memory** = "笔记本"(永久保存,随时翻阅)
|
||||
| 方式 | 比喻 | 特点 |
|
||||
|------|------|------|
|
||||
| 内存存储 | "草稿纸" | 进程退出就没了 |
|
||||
| Memory | "笔记本" | 永久保存,随时翻阅 |
|
||||
|
||||
---
|
||||
|
||||
## 关键概念
|
||||
|
||||
> **再次强调**:以下 Session、Store 等概念都是**业务层实现**,用于管理对话历史的存储。Eino 框架本身不提供这些组件,而是由业务层负责管理消息列表,然后将消息传递给 `adk.Runner` 进行处理。
|
||||
> [!tip] 重要提示
|
||||
> 以下 Session、Store 等概念都是**业务层实现**,用于管理对话历史的存储。Eino 框架本身不提供这些组件,而是由业务层负责管理消息列表,然后将消息传递给 `adk.Runner` 进行处理。
|
||||
|
||||
### Session(业务层概念)
|
||||
|
||||
@@ -129,10 +137,12 @@ type Store struct {
|
||||
|
||||
**为什么用 JSONL?**
|
||||
|
||||
- **简单**:每行一个 JSON 对象,易于读写
|
||||
- **可扩展**:可以追加新消息,无需重写整个文件
|
||||
- **可读性好**:可以用文本编辑器直接查看
|
||||
- **容错性强**:单行损坏不影响其他行
|
||||
| 特性 | 说明 |
|
||||
|------|------|
|
||||
| **简单** | 每行一个 JSON 对象,易于读写 |
|
||||
| **可扩展** | 可以追加新消息,无需重写整个文件 |
|
||||
| **可读性好** | 可以用文本编辑器直接查看 |
|
||||
| **容错性强** | 单行损坏不影响其他行 |
|
||||
|
||||
## Memory 的实现(业务层示例)
|
||||
|
||||
@@ -184,134 +194,137 @@ if err := session.Append(assistantMsg); err != nil {
|
||||
}
|
||||
```
|
||||
|
||||
**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch03/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch03/main.go)):
|
||||
### 关键代码解析
|
||||
|
||||
完整流程可以浓缩为以下核心片段:
|
||||
|
||||
```go
|
||||
// 创建或恢复 Session
|
||||
// 1. 创建或恢复 Session
|
||||
session, err := store.GetOrCreate(sessionID)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
// 用户输入
|
||||
// 2. 读取用户输入并追加到会话
|
||||
userMsg := schema.UserMessage(line)
|
||||
if err := session.Append(userMsg); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
// 调用 Agent
|
||||
// 3. 获取全部历史,送入 Agent 处理
|
||||
history := session.GetMessages()
|
||||
events := runner.Run(ctx, history)
|
||||
content := collectAssistantFromEvents(events)
|
||||
|
||||
// 保存助手回复
|
||||
// 4. 收集回复并存回会话
|
||||
assistantMsg := schema.AssistantMessage(content, nil)
|
||||
if err := session.Append(assistantMsg); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
```
|
||||
|
||||
> [!note] 简化说明
|
||||
> 以上代码已省略错误处理之外的细节(如命令行参数解析、事件流消费等),不能直接运行。完整代码请参考 [cmd/ch03/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch03/main.go)。
|
||||
|
||||
## Session 与 Agent 的关系:业务层与框架层的协作
|
||||
|
||||
**关键理解:**
|
||||
> [!tip] 理解要点
|
||||
> - **Session 是业务层概念**:由你的代码实现和管理,负责存储和加载对话历史
|
||||
> - **Agent(Runner)是框架层概念**:由 Eino 框架提供,负责处理消息并生成回复
|
||||
> - **两者的交互点**:业务层通过 `session.GetMessages()` 获取消息列表,传递给 `runner.Run(ctx, history)` 进行处理
|
||||
|
||||
- **Session 是业务层概念**:由业务代码实现和管理,负责存储和加载对话历史
|
||||
- **Agent(Runner)是框架层概念**:由 Eino 框架提供,负责处理消息并生成回复
|
||||
- **两者的交互点**:业务层通过 `session.GetMessages()` 获取消息列表,传递给 `runner.Run(ctx, history)` 进行处理
|
||||
**数据流示意图:**
|
||||
|
||||
**架构分层:**
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["用户输入"] --> B["session.Append()<br/>保存用户消息"]
|
||||
B --> C["session.GetMessages()<br/>获取完整历史"]
|
||||
C --> D["runner.Run(history)<br/>Agent 处理消息"]
|
||||
D --> E["收集助手回复"]
|
||||
E --> F["session.Append()<br/>保存助手消息"]
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────────────────────────┐
|
||||
│ 业务层(你的代码) │
|
||||
│ ┌─────────────┐ ┌──────────────┐ ┌───────────────┐ │
|
||||
│ │ Session │───→│ GetMessages() │───→│ runner.Run() │ │
|
||||
│ │ (存储) │ │ (消息列表) │ │ (框架调用) │ │
|
||||
│ └─────────────┘ └──────────────┘ └───────────────┘ │
|
||||
│ ↑ │ │
|
||||
│ │ ↓ │
|
||||
│ ┌─────────────┐ ┌───────────────┐ │
|
||||
│ │ Append() │←─────────────────────│ 助手回复 │ │
|
||||
│ │ (保存消息) │ └───────────────┘ │
|
||||
│ └─────────────┘ │
|
||||
└─────────────────────────────────────────────────────────────┘
|
||||
│
|
||||
↓
|
||||
┌─────────────────────────────────────────────────────────────┐
|
||||
│ 框架层(Eino 框架) │
|
||||
│ ┌───────────────────────────────────────────────────────┐ │
|
||||
│ │ adk.Runner:接收消息列表,调用 ChatModel,返回回复 │ │
|
||||
│ └───────────────────────────────────────────────────────┘ │
|
||||
└─────────────────────────────────────────────────────────────┘
|
||||
style B fill:#e8f5e9
|
||||
style C fill:#e8f5e9
|
||||
style F fill:#e8f5e9
|
||||
style D fill:#fff3e0
|
||||
```
|
||||
|
||||
**流程图:**
|
||||
**分层面貌:**
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 用户输入 │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ session.Append() │
|
||||
│ 保存用户消息 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ session.GetMessages()│
|
||||
│ 获取完整历史 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ runner.Run(history) │
|
||||
│ Agent 处理消息 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 收集助手回复 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ session.Append() │
|
||||
│ 保存助手消息 │
|
||||
└──────────────────────┘
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph biz_layer["业务层 — 你的代码"]
|
||||
S1["Session<br/>持久化存储"]
|
||||
S2["GetMessages()"]
|
||||
S3["Append()<br/>保存消息"]
|
||||
S1 --> S2
|
||||
S3 -.->|"回写"| S1
|
||||
end
|
||||
|
||||
subgraph frame_layer["框架层 — Eino"]
|
||||
R1["runner.Run()"]
|
||||
end
|
||||
|
||||
S2 --> R1
|
||||
R1 -->|"助手回复"| S3
|
||||
|
||||
style S1 fill:#e3f2fd
|
||||
style S2 fill:#e3f2fd
|
||||
style S3 fill:#e3f2fd
|
||||
style R1 fill:#fff3e0
|
||||
style biz_layer fill:none,stroke:#90caf9
|
||||
style frame_layer fill:none,stroke:#ffcc80
|
||||
```
|
||||
|
||||
## 本章小结
|
||||
|
||||
**框架层 vs 业务层:**
|
||||
| 核心概念 | 定位 | 说明 |
|
||||
|----------|------|------|
|
||||
| **Memory** | 业务层 | 对话历史的持久化存储,支持跨进程恢复 |
|
||||
| **Session** | 业务层 | 一次完整的对话会话,包含 ID、创建时间、消息列表 |
|
||||
| **Store** | 业务层 | 管理多个 Session 的存储,支持创建、获取、列表、删除 |
|
||||
| **JSONL 格式** | 业务层 | 简单的文件格式,易于读写和扩展 |
|
||||
| **adk.Runner** | 框架层 | 接收消息列表,调用 ChatModel,返回回复 |
|
||||
|
||||
- **Eino 框架层**:提供 `adk.Runner`、`schema.Message` 等基础抽象,不关心消息如何存储
|
||||
- **业务层(本章实现)**:Memory/Session/Store 是业务层概念,用于管理对话历史的存储
|
||||
> [!success] 学习成果
|
||||
> 完成本章后,你应该能够:
|
||||
> - 理解 Memory/Session/Store 的业务层职责
|
||||
> - 知道如何将对话历史持久化为 JSONL 文件
|
||||
> - 明白业务层与框架层的协作边界
|
||||
> - 根据业务需求选择合适的存储方案
|
||||
|
||||
**业务层概念:**
|
||||
|
||||
- **Memory**:对话历史的持久化存储,支持跨进程恢复
|
||||
- **Session**:一次完整的对话会话,包含 ID、创建时间、消息列表
|
||||
- **Store**:管理多个 Session 的存储,支持创建、获取、列表、删除
|
||||
- **JSONL 格式**:简单的文件格式,易于读写和扩展
|
||||
|
||||
**业务层与框架层的交互:**
|
||||
|
||||
- 业务层负责存储消息,通过 `session.GetMessages()` 获取消息列表
|
||||
- 将消息列表传递给框架层的 `runner.Run(ctx, history)` 进行处理
|
||||
- 收集框架层返回的回复,再由业务层保存到存储中
|
||||
|
||||
> **💡 提示**:本章的实现只是众多存储方案中的一种简单示例。在实际项目中,你可以根据业务需求选择数据库、Redis、云存储等方案,甚至可以实现更复杂的功能如会话过期清理、搜索、分享等。
|
||||
---
|
||||
|
||||
## 扩展思考:业务层存储方案的选择
|
||||
|
||||
本章提供的 JSONL 文件存储方案适合简单的单机应用。在实际业务中,你可能需要考虑其他存储方案:
|
||||
|
||||
**其他存储实现:**
|
||||
<table>
|
||||
<tr><td>存储方案</td><td>适用场景</td><td>优势</td><td>劣势</td></tr>
|
||||
<tr><td><strong>JSONL 文件</strong></td><td>单机应用、开发调试</td><td>零依赖,简单直观</td><td>不支持并发、分布式</td></tr>
|
||||
<tr><td><strong>SQLite / LevelDB</strong></td><td>桌面端应用</td><td>轻量级嵌入式数据库</td><td>不适合高并发写入</td></tr>
|
||||
<tr><td><strong>MySQL / PostgreSQL</strong></td><td>服务端部署</td><td>成熟稳定,功能丰富</td><td>运维成本较高</td></tr>
|
||||
<tr><td><strong>Redis</strong></td><td>分布式、高频访问</td><td>性能极高,支持过期策略</td><td>数据需额外持久化</td></tr>
|
||||
<tr><td><strong>S3 / OSS</strong></td><td>海量冷数据归档</td><td>成本极低,无限扩展</td><td>不适合频繁查询</td></tr>
|
||||
</table>
|
||||
|
||||
- 数据库存储(MySQL、PostgreSQL、MongoDB)
|
||||
- Redis 存储(支持分布式)
|
||||
- 云存储(S3、OSS)
|
||||
**高级功能展望:**
|
||||
|
||||
**高级功能:**
|
||||
- 会话过期清理(TTL 自动删除)
|
||||
- 会话全文搜索
|
||||
- 会话导出 / 导入
|
||||
- 会话分享(生成公开链接)
|
||||
|
||||
- 会话过期清理
|
||||
- 会话搜索
|
||||
- 会话导出/导入
|
||||
- 会话分享
|
||||
> [!tip] Middleware 联动
|
||||
> 当对话非常长时,单纯增加存储容量是不够的。第五章介绍的 **Summarization Middleware** 可以在调用 Agent 之前自动压缩历史消息,有效控制 Token 消耗。
|
||||
|
||||
---
|
||||
|
||||
## 下一章预告
|
||||
|
||||
[[Eino/quick_start/chapter_04_tool_and_filesystem|第四章:Tool 与文件系统]] 将为 Agent 添加文件访问能力,让智能助手能够读取代码仓库中的真实内容。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/chapter_02_chatmodelagent_runner_agentevent]]
|
||||
- [[Eino/quick_start/chapter_04_tool_and_filesystem]]
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-12"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第四章:Tool 与文件系统访问
|
||||
weight: 4
|
||||
tags: ["Eino", "Agent", "Tool", "Backend", "DeepAgent", "文件系统"]
|
||||
create time: "2026-04-29 15:30"
|
||||
---
|
||||
|
||||
本章目标:为 Agent 添加 Tool 能力,让 Agent 能够访问文件系统。
|
||||
# 第四章:Tool 与文件系统访问
|
||||
|
||||
## 概述
|
||||
|
||||
本章为 Agent 引入 Tool(工具)能力,使其能够突破纯文本对话的边界,直接操作文件系统、搜索代码库、执行命令。通过 DeepAgent 预构建组件和 Backend 抽象接口,只需几行配置即可让 Agent「看见」并「触碰」真实世界。
|
||||
|
||||
## 为什么需要 Tool
|
||||
|
||||
@@ -27,26 +27,37 @@ weight: 4
|
||||
|
||||
**简单类比:**
|
||||
|
||||
- **Agent** = "智能助手"(能理解指令,但需要工具才能执行)
|
||||
- **Tool** = "工具箱"(文件操作、网络请求、数据库查询等)
|
||||
- **Agent** = "智能助手"(能理解指令,但需要工具才能执行)
|
||||
- **Tool** = "工具箱"(文件操作、网络请求、数据库查询等)
|
||||
|
||||
> [!question] 深入思考
|
||||
>
|
||||
> 如果 Agent 拥有无限个 Tool,会不会反而变得更差?
|
||||
> 提示:考虑模型上下文窗口限制、Token 成本、以及"选择困难症"效应。实际设计中,**工具的元信息描述质量**比数量更重要——一个好的 `Description` 能让模型精准选对工具。
|
||||
|
||||
## 为什么需要文件系统能力
|
||||
|
||||
本示例是 ChatWithDoc(与文档对话),目标是帮助用户学习 Eino 框架并编写 Eino 代码。那么,最好的文档是什么?
|
||||
本示例是 ChatWithDoc(与文档对话),目标是帮助用户学习 Eino 框架并编写 Eino 代码。那么,最好的文档是什么?
|
||||
|
||||
**答案就是:Eino 仓库的代码本身。**
|
||||
**答案就是:Eino 仓库的代码本身。**
|
||||
|
||||
- **Code**: 源代码展示了框架的真实实现
|
||||
- **Comment**: 代码注释提供了设计思路和使用说明
|
||||
- **Examples**: 示例代码演示了最佳实践
|
||||
| 资源类型 | 价值 |
|
||||
|---------|------|
|
||||
| **Code** | 源代码展示了框架的真实实现 |
|
||||
| **Comment** | 代码注释提供了设计思路和使用说明 |
|
||||
| **Examples** | 示例代码演示了最佳实践 |
|
||||
|
||||
通过文件系统访问能力,Agent 可以直接读取 Eino 源码、注释和示例,为用户提供最准确、最及时的技术支持。
|
||||
通过文件系统访问能力,Agent 可以直接读取 Eino 源码、注释和示例,为用户提供最准确、最及时的技术支持。
|
||||
|
||||
> [!tip] 现实启发
|
||||
>
|
||||
> 很多官方文档会过时或写得模糊,但代码不会。让 Agent "读源码"是一种绕过信息衰减的可靠策略——这也是为什么 RAG 系统的向量库经常直接索引代码仓库的原因。
|
||||
|
||||
## 关键概念
|
||||
|
||||
### Tool 接口
|
||||
|
||||
`Tool` 是 Eino 中定义可执行能力的接口:
|
||||
`Tool` 是 Eino 中定义可执行能力的接口:
|
||||
|
||||
```go
|
||||
// BaseTool 提供工具的元信息,ChatModel 使用这些信息决定是否以及如何调用工具
|
||||
@@ -69,41 +80,51 @@ type StreamableTool interface {
|
||||
}
|
||||
```
|
||||
|
||||
**接口层次:**
|
||||
**接口层次:**
|
||||
|
||||
- `BaseTool`:基础接口,只提供元信息
|
||||
- `InvokableTool`:可执行工具(继承 BaseTool)
|
||||
- `StreamableTool`:流式工具(继承 BaseTool)
|
||||
- `BaseTool`:**信息层**,只提供工具的名称、描述和参数 schema,ChatModel 据此决定是否调用
|
||||
- `InvokableTool`:**同步执行层**,返回完整结果字符串,适合大多数文件操作场景
|
||||
- `StreamableTool`:**流式执行层**,通过 `StreamReader` 逐步输出结果,适合大量输出的场景(如长文件读取或命令执行)
|
||||
|
||||
> [!note] 设计要点
|
||||
>
|
||||
> 为什么 Tool 的参数是 JSON 字符串而不是 Go struct?
|
||||
> 因为 Tool 需要在 ChatModel 的上下文里传递——模型只能理解文本。所以 Eino 将参数序列化为 JSON 传给模型,模型再返回 JSON,由 ToolsNode 反序列化后调用底层实现。这种设计让任何语言/协议的工具都能与基于 JSON 的大模型对齐。
|
||||
|
||||
### Backend 接口
|
||||
|
||||
`Backend` 是 Eino 中用于文件系统操作的抽象接口:
|
||||
`Backend` 是 Eino 中用于文件系统操作的抽象接口,定义了六种核心文件能力:
|
||||
|
||||
```go
|
||||
type Backend interface {
|
||||
// 列出目录下的文件信息
|
||||
LsInfo(ctx context.Context, req *LsInfoRequest) ([]FileInfo, error)
|
||||
// 列出目录下的文件信息
|
||||
LsInfo(ctx context.Context, req *LsInfoRequest) ([]FileInfo, error)
|
||||
|
||||
// 读取文件内容,支持按行偏移和限制
|
||||
Read(ctx context.Context, req *ReadRequest) (*FileContent, error)
|
||||
// 读取文件内容,支持按行偏移和限制
|
||||
Read(ctx context.Context, req *ReadRequest) (*FileContent, error)
|
||||
|
||||
// 在文件中搜索匹配的内容
|
||||
GrepRaw(ctx context.Context, req *GrepRequest) ([]GrepMatch, error)
|
||||
// 在文件中搜索匹配的内容
|
||||
GrepRaw(ctx context.Context, req *GrepRequest) ([]GrepMatch, error)
|
||||
|
||||
// 根据 glob 模式匹配文件
|
||||
GlobInfo(ctx context.Context, req *GlobInfoRequest) ([]FileInfo, error)
|
||||
// 根据 glob 模式匹配文件
|
||||
GlobInfo(ctx context.Context, req *GlobInfoRequest) ([]FileInfo, error)
|
||||
|
||||
// 写入文件内容
|
||||
Write(ctx context.Context, req *WriteRequest) error
|
||||
// 写入文件内容
|
||||
Write(ctx context.Context, req *WriteRequest) error
|
||||
|
||||
// 编辑文件内容(字符串替换)
|
||||
Edit(ctx context.Context, req *EditRequest) error
|
||||
// 编辑文件内容(字符串替换)
|
||||
Edit(ctx context.Context, req *EditRequest) error
|
||||
}
|
||||
```
|
||||
|
||||
> [!tip] 方法分组
|
||||
>
|
||||
> Backend 的六类方法可以归为三类:
|
||||
> **发现**(LsInfo / GlobInfo)、**检索**(Read / GrepRaw)、**修改**(Write / Edit)。记住这个分类有助于快速理解不同 Backend 实现的侧重点。
|
||||
|
||||
### LocalBackend
|
||||
|
||||
`LocalBackend` 是 Backend 的本地文件系统实现,直接访问操作系统的文件系统:
|
||||
`LocalBackend` 是 Backend 的本地文件系统实现,直接访问操作系统的文件系统:
|
||||
|
||||
```go
|
||||
import localbk "github.com/cloudwego/eino-ext/adk/backend/local"
|
||||
@@ -111,13 +132,17 @@ import localbk "github.com/cloudwego/eino-ext/adk/backend/local"
|
||||
backend, err := localbk.NewBackend(ctx, &localbk.Config{})
|
||||
```
|
||||
|
||||
**特点:**
|
||||
**特点:**
|
||||
|
||||
- 直接访问本地文件系统,使用 Go 标准库实现
|
||||
- 直接访问本地文件系统,使用 Go 标准库实现
|
||||
- 支持所有 Backend 接口方法
|
||||
- 支持执行 shell 命令(ExecuteStreaming)
|
||||
- 路径安全:要求使用绝对路径,防止目录遍历攻击
|
||||
- 零配置:开箱即用,无需额外设置
|
||||
- 支持执行 shell 命令(ExecuteStreaming)
|
||||
- 路径安全:要求使用绝对路径,防止目录遍历攻击
|
||||
- 零配置:开箱即用,无需额外设置
|
||||
|
||||
> [!warning] 安全问题
|
||||
>
|
||||
> LocalBackend 要求使用绝对路径来防止 `../` 目录遍历攻击。这意味着 Agent 永远不能通过构造恶意文件名跳出设定的根目录范围——这是 LLM 集成中的关键安全措施。
|
||||
|
||||
## 实现:使用 DeepAgent
|
||||
|
||||
@@ -130,63 +155,75 @@ backend, err := localbk.NewBackend(ctx, &localbk.Config{})
|
||||
**ChatModelAgent vs DeepAgent 对比:**
|
||||
|
||||
<table>
|
||||
<tr><td>能力</td><td>ChatModelAgent</td><td>DeepAgent</td></tr>
|
||||
<tr><th>能力</th><th>ChatModelAgent</th><th>DeepAgent</th></tr>
|
||||
<tr><td>多轮对话</td><td>✅</td><td>✅</td></tr>
|
||||
<tr><td>添加自定义 Tool</td><td>✅ 手动注册每个 Tool</td><td>✅ 手动注册或自动注册</td></tr>
|
||||
<tr><td>文件系统访问(Backend)</td><td>❌ 需手动创建并注册所有文件工具</td><td>✅ 一级配置,自动注册</td></tr>
|
||||
<tr><td>命令执行(StreamingShell)</td><td>❌ 需手动创建</td><td>✅ 一级配置,自动注册</td></tr>
|
||||
<tr><td>内置任务管理</td><td>❌</td><td>✅ <pre>write_todos</pre> 工具</td></tr>
|
||||
<tr><td>内置任务管理</td><td>❌</td><td>✅ `write_todos` 工具</td></tr>
|
||||
<tr><td>支持子 Agent</td><td>❌</td><td>✅</td></tr>
|
||||
</table>
|
||||
|
||||
**选择建议:**
|
||||
|
||||
- 纯对话场景(无外部访问)→ 用 `ChatModelAgent`
|
||||
- 需要访问文件系统或执行命令 → 用 `DeepAgent`
|
||||
> [!tip] 选择建议
|
||||
>
|
||||
> - 纯对话场景(无外部访问)→ 用 `ChatModelAgent`
|
||||
> - 需要访问文件系统或执行命令 → 用 `DeepAgent`
|
||||
|
||||
### 为什么使用 DeepAgent?
|
||||
|
||||
相比直接使用 ChatModelAgent,DeepAgent 的优势:
|
||||
相比直接使用 ChatModelAgent,DeepAgent 的优势在于把「基础设施」封装成了第一类配置,你只需声明意图而非实现细节:
|
||||
|
||||
1. **一级配置**: Backend 和 StreamingShell 是一级配置,直接传入即可
|
||||
2. **自动注册工具**: 配置 Backend 后自动注册文件系统工具,无需手动创建
|
||||
3. **内置任务管理**: 提供 `write_todos` 工具,支持任务规划和跟踪
|
||||
4. **支持子 Agent**: 可以配置专门的子 Agent 处理特定任务
|
||||
5. **更强大**: 集成了文件系统、命令执行等多种能力
|
||||
1. **一级配置**:Backend 和 StreamingShell 直接在 Config 中传入,无需自己组装 ToolsNode
|
||||
2. **自动注册工具**:配置 Backend 后自动注册文件系统工具,免去逐个定义的样板代码
|
||||
3. **内置任务管理**:提供 `write_todos` 工具,支持复杂任务的规划与跟踪
|
||||
4. **支持子 Agent**:可以配置专门的子 Agent 处理特定任务
|
||||
5. **更强大**:集成了文件系统、命令执行等多种能力
|
||||
|
||||
### 代码实现
|
||||
|
||||
这段代码完成了两件事:创建 Backend 实例,再将其注入 DeepAgent。让我们逐行看:
|
||||
|
||||
```go
|
||||
import (
|
||||
localbk "github.com/cloudwego/eino-ext/adk/backend/local"
|
||||
"github.com/cloudwego/eino/adk/prebuilt/deep"
|
||||
localbk "github.com/cloudwego/eino-ext/adk/backend/local"
|
||||
"github.com/cloudwego/eino/adk/prebuilt/deep"
|
||||
)
|
||||
|
||||
// 创建 LocalBackend
|
||||
// 第一步:创建 LocalBackend —— 文件系统能力的实际执行者
|
||||
backend, err := localbk.NewBackend(ctx, &localbk.Config{})
|
||||
|
||||
// 创建 DeepAgent,自动注册文件系统工具
|
||||
// 第二步:将 backend 注入 DeepAgent,它会自动注册文件相关 Tool
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
Name: "Ch04ToolAgent",
|
||||
Description: "ChatWithDoc agent with filesystem access via LocalBackend.",
|
||||
ChatModel: cm,
|
||||
Instruction: instruction,
|
||||
Backend: backend, // 提供文件系统操作能力
|
||||
StreamingShell: backend, // 提供命令执行能力
|
||||
MaxIteration: 50,
|
||||
Name: "Ch04ToolAgent", // Agent 的名称
|
||||
Description: "ChatWithDoc agent with filesystem access.",
|
||||
ChatModel: cm, // 底层大模型
|
||||
Instruction: instruction, // 系统提示词
|
||||
Backend: backend, // 文件系统操作能力
|
||||
StreamingShell: backend, // 命令执行能力
|
||||
MaxIteration: 50, // 最大思考-行动循环次数
|
||||
})
|
||||
```
|
||||
|
||||
> [!note] MaxIteration 的含义
|
||||
>
|
||||
> 每次 Agent 「思考 → Tool Call → 观察结果」算一次迭代。设为 50 意味着最多允许 50 轮。如果达到上限仍未得到满意结果,会返回部分完成的内容。**设置过高会浪费 Token,过低则可能让 Agent 中途放弃**。后续章节会讨论如何优化这个值。
|
||||
|
||||
### DeepAgent 自动注册的工具
|
||||
|
||||
当配置了 `Backend` 和 `StreamingShell` 后,DeepAgent 会自动注册以下工具:
|
||||
当配置了 `Backend` 和 `StreamingShell` 后,DeepAgent 会自动注册以下工具:
|
||||
|
||||
- `read_file`: 读取文件内容
|
||||
- `write_file`: 写入文件内容
|
||||
- `edit_file`: 编辑文件内容
|
||||
- `glob`: 根据 glob 模式查找文件
|
||||
- `grep`: 在文件中搜索内容
|
||||
- `execute`: 执行 shell 命令
|
||||
| 工具名 | 对应 Backend 方法 | 用途 |
|
||||
|-------|------------------|------|
|
||||
| `read_file` | `Read` | 读取文件内容 |
|
||||
| `write_file` | `Write` | 写入文件内容 |
|
||||
| `edit_file` | `Edit` | 编辑文件内容 |
|
||||
| `glob` | `GlobInfo` | 根据 glob 模式查找文件 |
|
||||
| `grep` | `GrepRaw` | 在文件中搜索内容 |
|
||||
| `execute` | shell 命令 | 执行 shell 命令 |
|
||||
|
||||
> [!note] 从接口到 Tool 的映射
|
||||
>
|
||||
> Backend 定义了 6 个方法(LsInfo / Read / GrepRaw / GlobInfo / Write / Edit),但 DeepAgent 只注册了 5 个 Tool(少了 LsInfo)。这是因为 `glob` 已经能完成目录探索的需求。如果未来需要更详细的列表信息,可以手动补充 LsInfo 对应的 Tool。
|
||||
|
||||
## 代码位置
|
||||
|
||||
@@ -221,17 +258,27 @@ go run ./cmd/ch04
|
||||
|
||||
**推荐的三仓库目录结构(如要完整体验):**
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
ER[eino 核心库\nPROJECT_ROOT] --> ADK[adk/]
|
||||
ER --> COMP[components/]
|
||||
ER --> COMPOSE[compose/]
|
||||
ER --> EXT[eino-ext 扩展库]
|
||||
ER --> EX[eino-examples 示例库]
|
||||
EX --> QS[quickstart/]
|
||||
QS --> CW[chatwitheino / 本示例]
|
||||
|
||||
style ER fill:#e3f2fd
|
||||
style CW fill:#fff3e0
|
||||
```
|
||||
eino/ # PROJECT_ROOT(Eino 核心库)
|
||||
├── adk/
|
||||
├── components/
|
||||
├── compose/
|
||||
├── ext/ # eino-ext(扩展组件,如 OpenAI、Ark 等实现)
|
||||
├── examples/ # eino-examples(本仓库,本示例所在位置)
|
||||
│ └── quickstart/
|
||||
│ └── chatwitheino/
|
||||
└── ...
|
||||
```
|
||||
|
||||
> [!example] PROJECT_ROOT 的两种用法
|
||||
>
|
||||
> **场景一:快速验证**(不设置)
|
||||
> Agent 只能访问 `chatwitheino` 目录下的文件,适合学习当前示例的代码结构。
|
||||
>
|
||||
> **场景二:完整对话**(设置 `export PROJECT_ROOT=/path/to/eino`)
|
||||
> Agent 可以搜索 Eino 框架全部源码——比如查找某个 API 的实现细节、阅读组件设计思路等。这是 ChatWithEino 真正「与文档对话」的价值所在。
|
||||
|
||||
可以使用 `dev_setup.sh` 脚本自动设置上述目录结构:
|
||||
|
||||
@@ -264,66 +311,85 @@ you> 读取 main.go 文件的内容
|
||||
|
||||
## Tool 调用流程
|
||||
|
||||
当 Agent 需要调用 Tool 时:
|
||||
当 Agent 需要调用 Tool 时,内部会经历以下循环。你可以把整个过程理解为「思考 → 行动 → 观察」的闭环:
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
U[用户提问] --> A{Agent\n分析意图}
|
||||
A -->|纯对话| R[直接回复]
|
||||
A -->|需要工具| T1[生成 Tool Call\n参数 JSON]
|
||||
T1 --> T2[执行 Tool\n读取/写入/搜索等]
|
||||
T2 --> T3[返回 Tool Result]
|
||||
T3 --> A2{Agent\n整合信息}
|
||||
A2 --> R2[生成最终回复]
|
||||
R2 --> U2[回复用户]
|
||||
|
||||
style A fill:#e1f5fe
|
||||
style T2 fill:#fff3e0
|
||||
style T3 fill:#e8f5e9
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 用户:列出当前目录的文件 │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 分析意图 │
|
||||
│ 决定调用 glob 工具 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 生成 Tool Call │
|
||||
│ {"pattern": "*"} │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 执行 Tool │
|
||||
│ glob("*") │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 返回 Tool Result │
|
||||
│ {"files": [...]} │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 生成回复 │
|
||||
│ "找到 5 个文件..." │
|
||||
└──────────────────────┘
|
||||
```
|
||||
|
||||
**以"列出当前目录的文件"为例:**
|
||||
|
||||
> [!example] 逐步拆解
|
||||
>
|
||||
> **Step 1 — 意图识别**
|
||||
> 用户说"列出文件",Agent 判断这不是纯对话请求,而是文件系统操作意图。
|
||||
>
|
||||
> **Step 2 — 工具选择与参数生成**
|
||||
> Agent 从可用工具中选择 `glob`(查找文件列表),生成参数 `{"pattern": "*"}`。
|
||||
>
|
||||
> **Step 3 — Tool 执行**
|
||||
> `ToolsNode` 接收 JSON 参数,反序列化后调用 `GlobInfo`,将结果序列化为 JSON 返回。
|
||||
>
|
||||
> **Step 4 — 结果整合**
|
||||
> Agent 收到 `["main.go", "go.mod", ...]`,整理成自然语言回复:"找到 5 个文件..."。
|
||||
|
||||
你是否有想过——Agent 是如何知道该选哪个工具的?关键在于 `BaseTool.Info()` 返回的 **元信息**(名称、描述、参数 schema)。ChatModel 基于这些信息来决定调用哪个 Tool 以及传入什么参数。这就是后面会详细展开的 Function Calling 机制。
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Tool**:Agent 的能力扩展,让 Agent 能够执行具体操作
|
||||
- **Backend**:文件系统操作的抽象接口,提供统一的文件操作能力
|
||||
- **LocalBackend**:Backend 的本地文件系统实现,直接访问操作系统文件系统
|
||||
- **DeepAgent**:预构建的高级 Agent,提供 Backend 和 StreamingShell 的一级配置
|
||||
- **自动注册工具**:配置 Backend 后自动注册文件系统工具
|
||||
- **Tool 调用流程**:Agent 分析意图 → 生成 Tool Call → 执行 Tool → 返回结果 → 生成回复
|
||||
| 概念 | 一句话理解 |
|
||||
|------|-----------|
|
||||
| **Tool** | Agent 的能力扩展,让它能执行具体操作而非仅对话 |
|
||||
| **Backend** | 文件系统操作的抽象接口,定义发现、检索、修改三类方法 |
|
||||
| **LocalBackend** | Backend 的本地实现,开箱即用且路径安全 |
|
||||
| **DeepAgent** | 预构建的高级 Agent,配置 Backend 即可自动获得文件能力 |
|
||||
| **自动注册** | 声明式配置优于命令式组装——传一个 Backend 就得到五个工具 |
|
||||
| **调用流程** | 意图分析 → Tool Call → 执行 → 结果 → 回复(闭环迭代) |
|
||||
|
||||
## 扩展思考
|
||||
|
||||
**其他 Tool 类型:**
|
||||
### 其他 Tool 类型
|
||||
|
||||
- HTTP Tool:调用外部 API
|
||||
- Database Tool:查询数据库
|
||||
- Calculator Tool:执行计算
|
||||
- Code Executor Tool:运行代码
|
||||
当前我们只用了文件系统相关 Tool。Eino 生态中还支持更多类型:
|
||||
|
||||
**其他 Backend 实现:**
|
||||
- **HTTP Tool**:调用外部 API,让 Agent 具备联网能力
|
||||
- **Database Tool**:查询数据库,适用于数据分析场景
|
||||
- **Calculator Tool**:精确计算,弥补大模型不擅长数学的问题
|
||||
- **Code Executor Tool**:运行代码,适合生成并验证算法实现
|
||||
|
||||
- 可以基于 Backend 接口实现其他存储后端
|
||||
- 例如:云存储、数据库存储等
|
||||
- LocalBackend 已经提供了完整的文件系统操作能力
|
||||
### 自定义 Tool 创建
|
||||
|
||||
**自定义 Tool 创建:**
|
||||
除了使用 DeepAgent 自动注册的 Tool,你还可以手动创建。最简单的方式是使用 `utils.InferTool` 从函数签名自动推断参数 schema:
|
||||
|
||||
如果需要创建自定义 Tool,可以使用 `utils.InferTool` 从函数自动推断。详见:
|
||||
```go
|
||||
// 写一个普通 Go 函数...
|
||||
func Greet(name string) string {
|
||||
return fmt.Sprintf("Hello, %s!", name)
|
||||
}
|
||||
|
||||
// ...一行代码转成 Tool
|
||||
greetTool := utils.InferTool(Greet)
|
||||
```
|
||||
|
||||
如果你想了解更多,详见:
|
||||
|
||||
- [Tool 接口文档](https://github.com/cloudwego/eino/tree/main/components/tool)
|
||||
- [Tool 创建示例](https://github.com/cloudwego/eino-examples/tree/main/components/tool)
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/_index]]
|
||||
- [[Eino/quick_start/chapter_03_memory_and_session]] — Agent 的记忆与会话管理(上一章)
|
||||
- [[Eino/quick_start/chapter_05_middleware]] — Tool 调用中间件与拦截器(下一章)
|
||||
|
||||
@@ -1,65 +1,332 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-16"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第五章:Middleware(中间件模式)
|
||||
weight: 5
|
||||
tags: ["Eino", "Agent", "Middleware", "DeepAgent", "错误处理", "重试"]
|
||||
create time: "2026-04-29 16:00"
|
||||
---
|
||||
|
||||
本章目标:理解 Middleware 模式,实现 Tool 错误处理和 ChatModel 重试机制。
|
||||
# 第五章:Middleware(中间件模式)
|
||||
|
||||
## 概述
|
||||
|
||||
第四章为 Agent 加入了 Tool 能力后,Agent 已经可以「看见」和「触碰」真实世界了。但现实中的 API 会限流、文件会不存在、网络会超时——**直接暴露的错误会让 Agent 流程中断**。本章通过 Middleware 模式引入拦截器机制,让 Agent 具备错误自愈和自动重试的能力。
|
||||
|
||||
## 为什么需要 Middleware
|
||||
|
||||
第四章我们为 Agent 添加了 Tool 能力,让 Agent 能够访问文件系统。但在实际应用场景中,**Tool 报错或 ChatModel 报错是常见的现象**,例如:
|
||||
|
||||
- **Tool 报错**:文件不存在、参数错误、权限不足等
|
||||
- **ChatModel 报错**:API 限流(429)、网络超时、服务不可用等
|
||||
|
||||
### 问题一:Tool 错误会中断整个流程
|
||||
|
||||
当 Tool 执行失败时,错误会直接传播到 Agent,导致整个对话中断:
|
||||
第四章结束时,Tool 报错或 ChatModel 报错会直接中断整个对话流程:
|
||||
|
||||
```
|
||||
[tool call] read_file(file_path: "nonexistent.txt")
|
||||
Error: open nonexistent.txt: no such file or directory
|
||||
// 对话中断,用户需要重新开始
|
||||
// 💥 对话中断,用户需要重新开始
|
||||
```
|
||||
|
||||
### 问题二:模型调用可能因限流失败
|
||||
这类错误很常见:
|
||||
|
||||
当模型 API 返回 429(Too Many Requests)错误时,整个对话也会中断:
|
||||
| 场景 | 错误类型 | 常见原因 |
|
||||
|------|---------|---------|
|
||||
| **Tool 报错** | 业务错误 | 文件不存在、参数错误、权限不足 |
|
||||
| **ChatModel 报错** | 临时错误 | API 限流(429)、网络超时、服务不可用 |
|
||||
|
||||
```
|
||||
Error: rate limit exceeded (429)
|
||||
// 对话中断
|
||||
```
|
||||
> [!tip] 关键洞察
|
||||
>
|
||||
> 这些错误**不应该终止 Agent 流程**。更好的做法是把错误信息交给模型,让它自动调整策略继续执行:
|
||||
>
|
||||
> ```
|
||||
> [tool call] read_file(file_path: "nonexistent.txt")
|
||||
> [tool result] [tool error] open nonexistent.txt: no such file or directory
|
||||
> [assistant] 抱歉,文件不存在。让我先列出当前目录的文件...
|
||||
> [tool call] glob(pattern: "*")
|
||||
> // ✅ 对话继续,模型自行纠错
|
||||
> ```
|
||||
|
||||
### 期望的行为
|
||||
> [!question] 深入思考
|
||||
>
|
||||
> 既然可以直接把错误返回给模型,为什么不直接在每个 Tool 内部写 `if err != nil` 判断?
|
||||
> 提示:考虑开闭原则(OCP)——如果明天要加 10 个新 Tool,是不是每个都要改一遍?**Middleware 的本质是将横切关注点从业务代码中剥离**,这也是 AOP(面向切面编程)的核心思想。
|
||||
|
||||
这些报错信息往往**不希望直接终止 Agent 流程**,而是希望把报错信息给到模型,由模型自动纠错进行下一轮。例如:
|
||||
## 什么是 Middleware
|
||||
|
||||
```
|
||||
[tool call] read_file(file_path: "nonexistent.txt")
|
||||
[tool result] [tool error] open nonexistent.txt: no such file or directory
|
||||
[assistant] 抱歉,文件不存在。让我先列出当前目录的文件...
|
||||
[tool call] glob(pattern: "*")
|
||||
```
|
||||
**Middleware 是 Agent 的拦截器**,可以在调用前后插入自定义逻辑:
|
||||
|
||||
### Middleware 的定位
|
||||
|
||||
**Middleware 模式**可以扩展 Tool 和 ChatModel 的行为,非常适合解决这个问题:
|
||||
|
||||
- **Middleware 是 Agent 的拦截器**:在调用前后插入自定义逻辑
|
||||
- **Middleware 可处理错误**:将错误转换为模型可理解的格式
|
||||
- **Middleware 可实现重试**:自动重试失败的操作
|
||||
- **Middleware 可组合**:多个 Middleware 可以串联使用
|
||||
- **拦截调用**:在 Tool 或 ChatModel 执行前/后包装自定义行为
|
||||
- **错误转换**:将错误转为模型可理解的字符串,而非中断流程
|
||||
- **自动重试**:对临时错误(如限流)实现指数退避重试
|
||||
- **可组合**:多个 Middleware 串联形成责任链
|
||||
|
||||
**简单类比:**
|
||||
|
||||
- **Agent** = "业务逻辑"
|
||||
- **Middleware** = "AOP 切面"(日志、重试、错误处理等横切关注点)
|
||||
|
||||
> [!note] 装饰器模式
|
||||
>
|
||||
> Middleware 的本质是**装饰器模式**(Decorator Pattern)——每个 Middleware 包装原始调用,可以修改输入、输出或错误,而不改变被包装对象的接口。
|
||||
|
||||
## 核心概念
|
||||
|
||||
### Middleware 接口
|
||||
|
||||
`ChatModelAgentMiddleware` 是 Agent 中间件的统一接口:
|
||||
|
||||
```go
|
||||
type ChatModelAgentMiddleware interface {
|
||||
BeforeAgent(ctx context.Context, runCtx *ChatModelAgentContext) (context.Context, *ChatModelAgentContext, error)
|
||||
BeforeModelRewriteState(ctx context.Context, state *ChatModelAgentState, mc *ModelContext) (context.Context, *ChatModelAgentState, error)
|
||||
AfterModelRewriteState(ctx context.Context, state *ChatModelAgentState, mc *ModelContext) (context.Context, *ChatModelAgentState, error)
|
||||
WrapInvokableToolCall(ctx context.Context, endpoint InvokableToolCallEndpoint, tCtx *ToolContext) (InvokableToolCallEndpoint, error)
|
||||
WrapStreamableToolCall(ctx context.Context, endpoint StreamableToolCallEndpoint, tCtx *ToolContext) (StreamableToolCallEndpoint, error)
|
||||
WrapEnhancedInvokableToolCall(ctx context.Context, endpoint EnhancedInvokableToolCallEndpoint, tCtx *ToolContext) (EnhancedInvokableToolCallEndpoint, error)
|
||||
WrapEnhancedStreamableToolCall(ctx context.Context, endpoint EnhancedStreamableToolCallEndpoint, tCtx *ToolContext) (EnhancedStreamableToolCallEndpoint, error)
|
||||
WrapModel(ctx context.Context, m model.BaseChatModel, mc *ModelContext) (model.BaseChatModel, error)
|
||||
}
|
||||
```
|
||||
|
||||
**方法分组:**
|
||||
|
||||
| 分组 | 方法 | 作用时机 |
|
||||
|------|------|---------|
|
||||
| **Agent 生命周期** | `BeforeAgent` | 每次 Agent 运行前,可修改指令和工具配置 |
|
||||
| **状态处理** | `BeforeModelRewriteState` / `AfterModelRewriteState` | 每次模型调用前后的状态变换 |
|
||||
| **Tool 调用** | `WrapInvokableToolCall` / `WrapStreamableToolCall` | 包装同步/流式 Tool 的执行 |
|
||||
| **模型调用** | `WrapModel` | 包装底层 ChatModel 的调用 |
|
||||
|
||||
### 洋葱模型:Middleware 执行顺序
|
||||
|
||||
Handlers 按**数组正序**包装,形成洋葱模型:
|
||||
|
||||
```go
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&middlewareA{}, // 最外层:最先 Wrap,最后生效
|
||||
&middlewareB{}, // 中间层
|
||||
&middlewareC{}, // 最内层:最后 Wrap,最先生效
|
||||
}
|
||||
```
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
subgraph Request ["📥 请求方向 →"]
|
||||
A["Middleware A\n(最外层)"] --> B["Middleware B\n(中间层)"]
|
||||
B --> C["Middleware C\n(最内层)"]
|
||||
C --> T["实际 Tool/Model\n执行"]
|
||||
end
|
||||
|
||||
subgraph Response ["📤 响应方向 ←"]
|
||||
T --> CR["Middleware C\n返回"]
|
||||
CR --> CB["Middleware B\n返回"]
|
||||
CB --> CA["Middleware A\n返回"]
|
||||
end
|
||||
|
||||
style A fill:#fce4ec
|
||||
style B fill:#e8f5e9
|
||||
style C fill:#e3f2fd
|
||||
style T fill:#fff3e0
|
||||
```
|
||||
|
||||
> [!warning] 实用建议
|
||||
>
|
||||
> 将 `safeToolMiddleware`(错误捕获)放在最内层(数组末尾),确保其他 Middleware 抛出的中断错误能正确向外传播,不被吞掉。
|
||||
|
||||
### ModelRetryConfig:内置重试配置
|
||||
|
||||
`ModelRetryConfig` 提供了 ChatModel 级别的自动重试能力:
|
||||
|
||||
```go
|
||||
type ModelRetryConfig struct {
|
||||
MaxRetries int // 最大重试次数
|
||||
IsRetryAble func(ctx context.Context, err error) bool // 哪些错误可重试
|
||||
}
|
||||
```
|
||||
|
||||
**重试策略:**
|
||||
|
||||
| 策略 | 说明 |
|
||||
|------|------|
|
||||
| **指数退避** | 每次重试间隔递增,避免频繁请求加剧限流 |
|
||||
| **条件过滤** | 通过 `IsRetryAble` 精确控制哪些错误值得重试 |
|
||||
| **自动恢复** | 无需用户干预,模型调用失败后自动重试 |
|
||||
|
||||
## 实现细节
|
||||
|
||||
### SafeToolMiddleware:错误转换
|
||||
|
||||
`SafeToolMiddleware` 捕获 Tool 执行时的错误,将其转换为字符串返回给模型而非中断流程:
|
||||
|
||||
```go
|
||||
type safeToolMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *safeToolMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
result, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
// ❗ 中断错误不转换,需要继续向外传播
|
||||
if _, ok := compose.IsInterruptRerunError(err); ok {
|
||||
return "", err
|
||||
}
|
||||
// ✅ 普通错误转为字符串,交给模型处理
|
||||
return fmt.Sprintf("[tool error] %v", err), nil
|
||||
}
|
||||
return result, nil
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
**设计要点:**
|
||||
|
||||
- **区分错误类型**:中断错误(如主动要求停止)必须传播,业务错误(如文件不存在)可以转换
|
||||
- **不吞错**:只转换预期的业务错误,真正的系统异常仍向上抛出
|
||||
- **格式化**:使用 `[tool error]` 前缀方便模型识别并回复时引用
|
||||
|
||||
流式 Tool 的错误处理同理,需将错误封装为单帧流:
|
||||
|
||||
```go
|
||||
func (m *safeToolMiddleware) WrapStreamableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.StreamableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.StreamableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (*schema.StreamReader[string], error) {
|
||||
sr, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
if _, ok := compose.IsInterruptRerunError(err); ok {
|
||||
return nil, err
|
||||
}
|
||||
// 返回包含错误信息的单帧流
|
||||
return singleChunkReader(fmt.Sprintf("[tool error] %v", err)), nil
|
||||
}
|
||||
return safeWrapReader(sr), nil
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
### 注册 Middleware 与重试配置
|
||||
|
||||
将 Middleware 注入 DeepAgent 的配置中:
|
||||
|
||||
```go
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
Name: "Ch05MiddlewareAgent",
|
||||
Description: "ChatWithDoc agent with safe tool middleware and retry.",
|
||||
ChatModel: cm,
|
||||
Instruction: agentInstruction,
|
||||
Backend: backend,
|
||||
StreamingShell: backend,
|
||||
MaxIteration: 50,
|
||||
|
||||
// ⭐ 注册 Middleware
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&safeToolMiddleware{}, // 将 Tool 错误转为字符串
|
||||
},
|
||||
|
||||
// ⭐ 注册模型重试配置
|
||||
ModelRetryConfig: &adk.ModelRetryConfig{
|
||||
MaxRetries: 5,
|
||||
IsRetryAble: func(_ context.Context, err error) bool {
|
||||
return strings.Contains(err.Error(), "429") ||
|
||||
strings.Contains(err.Error(), "Too Many Requests")
|
||||
},
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
> [!note] Handlers vs Middlewares
|
||||
>
|
||||
> `Handlers` 字段(在 Config 中)和 "Middleware"(文档讨论的概念)是同一回事——`Handlers` 是配置字段名,而 `ChatModelAgentMiddleware` 是对接口的命名。
|
||||
|
||||
## 执行流程
|
||||
|
||||
结合 Middleware 后,一次 Tool 调用的完整生命周期如下:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
U["用户:读取不存在的文件"] --> A{"Agent 分析意图"}
|
||||
A -->|"决定调用 Tool"| M["SafeToolMiddleware\n拦截 Tool 调用"]
|
||||
M --> T["执行 read_file\n返回错误"]
|
||||
T --> E["SafeToolMiddleware\n捕获错误"]
|
||||
E -->|"非中断错误"| S["转换为字符串\ntool error: no such file"]
|
||||
E -->|"中断错误"| EP["向上抛出中断"]
|
||||
S --> R["返回 Tool Result"]
|
||||
R --> AG{"Agent 整合信息"}
|
||||
AG -->|"生成解释性回复"| O["抱歉,文件不存在...\n尝试列出目录"]
|
||||
AG -->|"需要更多信息"| A
|
||||
|
||||
style M fill:#e8f5e9
|
||||
style E fill:#fff3e0
|
||||
style S fill:#e3f2fd
|
||||
style EP fill:#ffebee
|
||||
```
|
||||
|
||||
> [!example] 逐步拆解
|
||||
>
|
||||
> **Step 1 — 用户输入**
|
||||
> 用户请求读取一个不存在的文件。
|
||||
>
|
||||
> **Step 2 — 意图分析**
|
||||
> Agent 判断需要文件系统操作,决定调用 `read_file` Tool。
|
||||
>
|
||||
> **Step 3 — Middleware 拦截**
|
||||
> `SafeToolMiddleware.WrapInvokableToolCall` 在 Tool 执行前被触发,注册了自己的回调逻辑。
|
||||
>
|
||||
> **Step 4 — Tool 执行**
|
||||
> 实际文件读取操作失败,返回 `open nonexistent.txt: no such file` 错误。
|
||||
>
|
||||
> **Step 5 — 错误转换**
|
||||
> Middleware 发现这不是中断错误,将其包装为 `[tool error] open nonexistent.txt: ...` 字符串。
|
||||
>
|
||||
> **Step 6 — Agent 自愈**
|
||||
> Agent 收到带错误的 Tool Result,理解后回复用户并调整策略(如改用 `glob` 列出可用文件)。
|
||||
|
||||
## 扩展:Eino 内置 Middleware
|
||||
|
||||
Eino 生态还提供了以下开箱即用的中间件:
|
||||
|
||||
<table>
|
||||
<tr><th>Middleware</th><th>功能说明</th></tr>
|
||||
<tr><td><strong>reduction</strong></td><td>工具输出缩减——当工具返回过长时自动截断并存入文件系统,防止上下文溢出</td></tr>
|
||||
<tr><td><strong>summarization</strong></td><td>对话历史摘要——Token 超阈值时自动生成摘要压缩历史,节省上下文空间</td></tr>
|
||||
<tr><td><strong>skill</strong></td><td>技能加载——让 Agent 按需动态加载预定义的 SKILL.md 知识包</td></tr>
|
||||
</table>
|
||||
|
||||
### 多 Middleware 组合示例
|
||||
|
||||
```go
|
||||
import (
|
||||
"github.com/cloudwego/eino/adk/middlewares/reduction"
|
||||
"github.com/cloudwego/eino/adk/middlewares/summarization"
|
||||
)
|
||||
|
||||
// 创建 reduction:管理工具输出长度
|
||||
reductionMW, _ := reduction.New(ctx, &reduction.Config{
|
||||
Backend: filesystemBackend,
|
||||
MaxLengthForTrunc: 50000,
|
||||
MaxTokensForClear: 30000,
|
||||
})
|
||||
|
||||
// 创建 summarization:自动压缩对话历史
|
||||
summarizationMW, _ := summarization.New(ctx, &summarization.Config{
|
||||
Model: chatModel,
|
||||
Trigger: &summarization.TriggerCondition{
|
||||
ContextTokens: 190000,
|
||||
},
|
||||
})
|
||||
|
||||
// 组合使用
|
||||
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
summarizationMW, // 外层:对话历史摘要
|
||||
reductionMW, // 内层:工具输出缩减
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
> [!question] 扩展思考
|
||||
>
|
||||
> 在这个例子中,`summarizationMW` 在外层、`reductionMW` 在内层。如果把顺序反过来,会有什么影响?试着根据洋葱模型的执行顺序推导一下。
|
||||
|
||||
## 代码位置
|
||||
|
||||
- 入口代码:[cmd/ch05/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch05/main.go)
|
||||
@@ -77,13 +344,11 @@ export PROJECT_ROOT=/path/to/eino # Eino 核心库根目录
|
||||
在 `examples/quickstart/chatwitheino` 目录下执行:
|
||||
|
||||
```bash
|
||||
# 设置项目根目录
|
||||
export PROJECT_ROOT=/path/to/your/project
|
||||
|
||||
go run ./cmd/ch05
|
||||
```
|
||||
|
||||
输出示例:
|
||||
**输出示例:**
|
||||
|
||||
```
|
||||
you> 列出当前目录的文件
|
||||
@@ -97,352 +362,19 @@ you> 读取一个不存在的文件
|
||||
[assistant] 抱歉,文件不存在...
|
||||
```
|
||||
|
||||
## 关键概念
|
||||
|
||||
### Middleware 接口
|
||||
|
||||
`ChatModelAgentMiddleware` 是 Agent 的中间件接口:
|
||||
|
||||
```go
|
||||
type ChatModelAgentMiddleware interface {
|
||||
// BeforeAgent is called before each agent run, allowing modification of
|
||||
// the agent's instruction and tools configuration.
|
||||
BeforeAgent(ctx context.Context, runCtx *ChatModelAgentContext) (context.Context, *ChatModelAgentContext, error)
|
||||
|
||||
// BeforeModelRewriteState is called before each model invocation.
|
||||
// The returned state is persisted to the agent's internal state and passed to the model.
|
||||
BeforeModelRewriteState(ctx context.Context, state *ChatModelAgentState, mc *ModelContext) (context.Context, *ChatModelAgentState, error)
|
||||
|
||||
// AfterModelRewriteState is called after each model invocation.
|
||||
// The input state includes the model's response as the last message.
|
||||
AfterModelRewriteState(ctx context.Context, state *ChatModelAgentState, mc *ModelContext) (context.Context, *ChatModelAgentState, error)
|
||||
|
||||
// WrapInvokableToolCall wraps a tool's synchronous execution with custom behavior.
|
||||
// This method is only called for tools that implement InvokableTool.
|
||||
WrapInvokableToolCall(ctx context.Context, endpoint InvokableToolCallEndpoint, tCtx *ToolContext) (InvokableToolCallEndpoint, error)
|
||||
|
||||
// WrapStreamableToolCall wraps a tool's streaming execution with custom behavior.
|
||||
// This method is only called for tools that implement StreamableTool.
|
||||
WrapStreamableToolCall(ctx context.Context, endpoint StreamableToolCallEndpoint, tCtx *ToolContext) (StreamableToolCallEndpoint, error)
|
||||
|
||||
// WrapEnhancedInvokableToolCall wraps an enhanced tool's synchronous execution.
|
||||
// This method is only called for tools that implement EnhancedInvokableTool.
|
||||
WrapEnhancedInvokableToolCall(ctx context.Context, endpoint EnhancedInvokableToolCallEndpoint, tCtx *ToolContext) (EnhancedInvokableToolCallEndpoint, error)
|
||||
|
||||
// WrapEnhancedStreamableToolCall wraps an enhanced tool's streaming execution.
|
||||
// This method is only called for tools that implement EnhancedStreamableTool.
|
||||
WrapEnhancedStreamableToolCall(ctx context.Context, endpoint EnhancedStreamableToolCallEndpoint, tCtx *ToolContext) (EnhancedStreamableToolCallEndpoint, error)
|
||||
|
||||
// WrapModel wraps a chat model with custom behavior.
|
||||
// This method is called at request time when the model is about to be invoked.
|
||||
WrapModel(ctx context.Context, m model.BaseChatModel, mc *ModelContext) (model.BaseChatModel, error)
|
||||
}
|
||||
```
|
||||
|
||||
**设计理念:**
|
||||
|
||||
- **装饰器模式**:每个 Middleware 包装原始调用,可以修改输入、输出或错误
|
||||
- **洋葱模型**:请求从外向内穿过 Middleware,响应从内向外返回
|
||||
- **可组合**:多个 Middleware 按顺序执行
|
||||
|
||||
### Middleware 执行顺序
|
||||
|
||||
`Handlers`(即 Middlewares)按**数组正序**包装,形成洋葱模型:
|
||||
|
||||
```go
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&middlewareA{}, // 最外层:最先 Wrap,最先拦截请求,但 WrapModel 最后生效
|
||||
&middlewareB{}, // 中间层
|
||||
&middlewareC{}, // 最内层:最后 Wrap
|
||||
}
|
||||
```
|
||||
|
||||
**对于 Tool 调用的执行顺序:**
|
||||
|
||||
```
|
||||
请求 → A.Wrap → B.Wrap → C.Wrap → 实际 Tool 执行 → C返回 → B返回 → A返回 → 响应
|
||||
```
|
||||
|
||||
**实用建议:** 将 `safeToolMiddleware`(错误捕获)放在最内层(数组末尾),确保其他 Middleware 抛出的中断错误能正确向外传播。
|
||||
|
||||
### SafeToolMiddleware
|
||||
|
||||
`SafeToolMiddleware` 将 Tool 错误转换为字符串,让模型能够理解并处理:
|
||||
|
||||
```go
|
||||
type safeToolMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *safeToolMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
result, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
// 将错误转换为字符串,而不是返回错误
|
||||
return fmt.Sprintf("[tool error] %v", err), nil
|
||||
}
|
||||
return result, nil
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
**效果:**
|
||||
|
||||
```
|
||||
[tool call] read_file(file_path: "nonexistent.txt")
|
||||
[tool result] [tool error] open nonexistent.txt: no such file or directory
|
||||
[assistant] 抱歉,文件不存在,请检查文件路径...
|
||||
// 对话继续,模型可以根据错误信息调整策略
|
||||
```
|
||||
|
||||
### ModelRetryConfig
|
||||
|
||||
`ModelRetryConfig` 配置 ChatModel 的自动重试:
|
||||
|
||||
```go
|
||||
type ModelRetryConfig struct {
|
||||
MaxRetries int // 最大重试次数
|
||||
IsRetryAble func(ctx context.Context, err error) bool // 判断是否可重试
|
||||
}
|
||||
```
|
||||
|
||||
**使用方式(以 DeepAgent 为例):**
|
||||
|
||||
```go
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
// ...
|
||||
ModelRetryConfig: &adk.ModelRetryConfig{
|
||||
MaxRetries: 5,
|
||||
IsRetryAble: func(_ context.Context, err error) bool {
|
||||
// 429 限流错误可重试
|
||||
return strings.Contains(err.Error(), "429") ||
|
||||
strings.Contains(err.Error(), "Too Many Requests") ||
|
||||
strings.Contains(err.Error(), "qpm limit")
|
||||
},
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
**重试策略:**
|
||||
|
||||
- 指数退避:每次重试间隔递增
|
||||
- 可配置条件:通过 `IsRetryAble` 判断哪些错误可重试
|
||||
- 自动恢复:无需用户干预
|
||||
|
||||
## Middleware 的实现
|
||||
|
||||
### 1. 实现 SafeToolMiddleware
|
||||
|
||||
```go
|
||||
type safeToolMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *safeToolMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
result, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
// 中断错误不转换,需要继续传播
|
||||
if _, ok := compose.IsInterruptRerunError(err); ok {
|
||||
return "", err
|
||||
}
|
||||
// 其他错误转换为字符串
|
||||
return fmt.Sprintf("[tool error] %v", err), nil
|
||||
}
|
||||
return result, nil
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
### 2. 实现流式 Tool 错误处理
|
||||
|
||||
```go
|
||||
func (m *safeToolMiddleware) WrapStreamableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.StreamableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.StreamableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (*schema.StreamReader[string], error) {
|
||||
sr, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
if _, ok := compose.IsInterruptRerunError(err); ok {
|
||||
return nil, err
|
||||
}
|
||||
// 返回包含错误信息的单帧流
|
||||
return singleChunkReader(fmt.Sprintf("[tool error] %v", err)), nil
|
||||
}
|
||||
// 包装流,捕获流中的错误
|
||||
return safeWrapReader(sr), nil
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
### 3. 配置 Agent 使用 Middleware
|
||||
|
||||
本章继续使用第四章引入的 `DeepAgent`,在其 `Handlers` 字段中注册 Middleware:
|
||||
|
||||
```go
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
Name: "Ch05MiddlewareAgent",
|
||||
Description: "ChatWithDoc agent with safe tool middleware and retry.",
|
||||
ChatModel: cm,
|
||||
Instruction: agentInstruction,
|
||||
Backend: backend,
|
||||
StreamingShell: backend,
|
||||
MaxIteration: 50,
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&safeToolMiddleware{}, // 将 Tool 错误转换为字符串
|
||||
},
|
||||
ModelRetryConfig: &adk.ModelRetryConfig{
|
||||
MaxRetries: 5,
|
||||
IsRetryAble: func(_ context.Context, err error) bool {
|
||||
return strings.Contains(err.Error(), "429") ||
|
||||
strings.Contains(err.Error(), "Too Many Requests")
|
||||
},
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
**注意**:`Handlers` 字段(在配置中)和 "Middleware"(在文档中讨论的概念)是同一回事——`Handlers` 是配置字段名,而 `ChatModelAgentMiddleware` 是接口名。
|
||||
|
||||
```
|
||||
**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch05/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch05/main.go)):
|
||||
|
||||
```go
|
||||
// SafeToolMiddleware 捕获 Tool 错误并转换为字符串
|
||||
type safeToolMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *safeToolMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
_ *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
result, err := endpoint(ctx, args, opts...)
|
||||
if err != nil {
|
||||
if _, ok := compose.IsInterruptRerunError(err); ok {
|
||||
return "", err
|
||||
}
|
||||
return fmt.Sprintf("[tool error] %v", err), nil
|
||||
}
|
||||
return result, nil
|
||||
}, nil
|
||||
}
|
||||
|
||||
// 配置 DeepAgent(与第四章一样,新增 Handlers 和 ModelRetryConfig)
|
||||
agent, _ := deep.New(ctx, &deep.Config{
|
||||
ChatModel: cm,
|
||||
Backend: backend,
|
||||
StreamingShell: backend,
|
||||
MaxIteration: 50,
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&safeToolMiddleware{},
|
||||
},
|
||||
ModelRetryConfig: &adk.ModelRetryConfig{
|
||||
MaxRetries: 5,
|
||||
IsRetryAble: func(_ context.Context, err error) bool {
|
||||
return strings.Contains(err.Error(), "429")
|
||||
},
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
## Middleware 执行流程
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 用户:读取不存在的文件 │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 分析意图 │
|
||||
│ 决定调用 read_file │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ SafeToolMiddleware │
|
||||
│ 拦截 Tool 调用 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 执行 read_file │
|
||||
│ 返回错误 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ SafeToolMiddleware │
|
||||
│ 将错误转换为字符串 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 返回 Tool Result │
|
||||
│ "[tool error] ..." │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 生成回复 │
|
||||
│ "抱歉,文件不存在..." │
|
||||
└──────────────────────┘
|
||||
```
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Middleware**:Agent 的拦截器,可以在调用前后插入自定义逻辑
|
||||
- **SafeToolMiddleware**:将 Tool 错误转换为字符串,让模型能够理解并处理
|
||||
- **ModelRetryConfig**:配置 ChatModel 的自动重试,处理限流等临时错误
|
||||
- **装饰器模式**:Middleware 包装原始调用,可以修改输入、输出或错误
|
||||
- **洋葱模型**:请求从外向内穿过 Middleware,响应从内向外返回
|
||||
| 概念 | 一句话理解 |
|
||||
|------|-----------|
|
||||
| **Middleware** | Agent 的拦截器,在调用前后插入自定义逻辑 |
|
||||
| **SafeToolMiddleware** | 将 Tool 错误转为字符串交给模型,而非中断流程 |
|
||||
| **ModelRetryConfig** | 配置 ChatModel 的自动重试,处理限流等临时错误 |
|
||||
| **洋葱模型** | 请求从外向内穿过 Middleware,响应从内向外返回 |
|
||||
| **装饰器模式** | 每个 Middleware 包装原始调用,可修改输入、输出或错误 |
|
||||
| **中断错误不转换** | 只有业务错误才转字符串,中断错误继续传播 |
|
||||
|
||||
## 扩展思考
|
||||
## 关联笔记
|
||||
|
||||
**Eino 内置 Middleware:**
|
||||
|
||||
<table>
|
||||
<tr><td>Middleware</td><td>功能说明</td></tr>
|
||||
<tr><td><strong>reduction</strong></td><td>工具输出缩减,当工具返回内容过长时自动截断并卸载到文件系统,防止上下文溢出</td></tr>
|
||||
<tr><td><strong>summarization</strong></td><td>对话历史自动摘要,当 token 数量超过阈值时自动生成摘要压缩历史</td></tr>
|
||||
<tr><td><strong>skill</strong></td><td>技能加载中间件,让 Agent 能够动态加载和执行预定义的技能</td></tr>
|
||||
</table>
|
||||
|
||||
**Middleware 链示例:**
|
||||
|
||||
```go
|
||||
import (
|
||||
"github.com/cloudwego/eino/adk/middlewares/reduction"
|
||||
"github.com/cloudwego/eino/adk/middlewares/summarization"
|
||||
"github.com/cloudwego/eino/adk/middlewares/skill"
|
||||
)
|
||||
|
||||
// 创建 reduction middleware:管理工具输出长度
|
||||
reductionMW, _ := reduction.New(ctx, &reduction.Config{
|
||||
Backend: filesystemBackend, // 存储后端
|
||||
MaxLengthForTrunc: 50000, // 单次工具输出最大长度
|
||||
MaxTokensForClear: 30000, // 触发清理的 token 阈值
|
||||
})
|
||||
|
||||
// 创建 summarization middleware:自动压缩对话历史
|
||||
summarizationMW, _ := summarization.New(ctx, &summarization.Config{
|
||||
Model: chatModel, // 用于生成摘要的模型
|
||||
Trigger: &summarization.TriggerCondition{
|
||||
ContextTokens: 190000, // 触发摘要的 token 阈值
|
||||
},
|
||||
})
|
||||
|
||||
// 组合多个 middleware(概念示例,使用 DeepAgent 时将 adk.NewChatModelAgent 替换为 deep.New)
|
||||
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
|
||||
Handlers: []adk.ChatModelAgentMiddleware{ // 注意:配置字段名为 Handlers,概念上与 Middlewares 等价
|
||||
summarizationMW, // 最外层:对话历史摘要
|
||||
reductionMW, // 中间层:工具输出缩减
|
||||
},
|
||||
})
|
||||
```
|
||||
- [[Eino/quick_start/_index]]
|
||||
- [[Eino/quick_start/chapter_04_tool_and_filesystem]] — Tool 与文件系统访问(上一章)
|
||||
- [[Eino/quick_start/chapter_06_callback_and_trace]] — Callback 与 Trace 可观测性(下一章)
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-12"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第六章:Callback 与 Trace(可观测性)
|
||||
weight: 6
|
||||
tags: [Eino, Callback, Trace, 可观测性, CozeLoop]
|
||||
create time: 2026-04-29 15:30
|
||||
---
|
||||
|
||||
本章目标:理解 Callback 机制,集成 CozeLoop 实现链路追踪和可观测性。
|
||||
# Eino 快速入门 · 第六章:Callback 与 Trace(可观测性)
|
||||
|
||||
## 概述
|
||||
|
||||
在构建 Agent 应用时,我们常常面临一个核心问题:**Agent 内部到底发生了什么?** 本章将介绍 Eino 的 Callback 机制——一套非侵入式的旁路钩子系统,让你能在不改动业务代码的前提下,获取组件生命周期的每一个关键信息。通过 Callback 配合 CozeLoop,你将获得完整的链路追踪、性能指标和错误定位能力。
|
||||
|
||||
## 代码位置
|
||||
|
||||
@@ -64,14 +64,10 @@ you> 你好
|
||||
- 不知道 Token 消耗了多少
|
||||
- 出问题时难以定位原因
|
||||
|
||||
**Callback 的定位:**
|
||||
> [!NOTE] Callback 定位
|
||||
> Callback 是 Eino 的**旁路机制**——从 component 到 compose,一以贯之。它在固定点位触发,可抽取实时信息(输入、输出、错误、流式数据),用途覆盖观测、日志、指标、追踪、调试、审计等场景。
|
||||
|
||||
- **Callback 是 Eino 的旁路机制**:从 component 到 compose(下文详谈)到 adk,一以贯之
|
||||
- **Callback 在固定点位触发**:组件生命周期的 5 个关键时机
|
||||
- **Callback 可抽取实时信息**:输入、输出、错误、流式数据等
|
||||
- **Callback 用途广泛**:观测、日志、指标、追踪、调试、审计等
|
||||
|
||||
**简单类比:**
|
||||
**类比理解:**
|
||||
|
||||
- **Agent** = "业务逻辑"(主路)
|
||||
- **Callback** = "旁路钩子"(在固定点位抽取信息)
|
||||
@@ -110,21 +106,21 @@ type Handler interface {
|
||||
- **状态传递**:同一 Handler 的 OnStart→OnEnd 可通过 context 传递状态
|
||||
- **性能优化**:实现 `TimingChecker` 接口可跳过不需要的时机
|
||||
|
||||
**RunInfo 结构:**
|
||||
> [!TIP] RunInfo 结构
|
||||
> `RunInfo` 携带了组件运行时身份,是日志和追踪中最重要的标识信息。
|
||||
|
||||
```go
|
||||
type RunInfo struct {
|
||||
Name string // 业务名称(节点名或用户指定)
|
||||
Type string // 实现类型(如 "OpenAI")
|
||||
Component string // 组件类型(如 "ChatModel")
|
||||
Name string // 业务名称(节点名或用户指定)
|
||||
Type string // 实现类型(如 "OpenAI")
|
||||
Component string // 组件类型(如 "ChatModel")
|
||||
}
|
||||
```
|
||||
|
||||
**重要提示:**
|
||||
|
||||
- 流式回调必须关闭 StreamReader,否则会导致 goroutine 泄漏
|
||||
- 不要修改 Input/Output,它们被所有下游共享
|
||||
- RunInfo 可能为 nil,使用前需要检查
|
||||
> [!IMPORTANT] 流式回调注意事项
|
||||
> - 流式回调必须关闭 StreamReader,否则会导致 goroutine 泄漏
|
||||
> - 不要修改 Input/Output,它们被所有下游共享
|
||||
> - RunInfo 可能为 nil,使用前需要检查
|
||||
|
||||
### CozeLoop
|
||||
|
||||
@@ -160,68 +156,70 @@ Callback 在组件生命周期的 5 个关键时机触发。下表中 `Timing*`
|
||||
|
||||
<table>
|
||||
<tr><td>时机常量</td><td>对应 Handler 方法</td><td>触发点</td><td>输入/输出</td></tr>
|
||||
<tr><td><pre>TimingOnStart</pre></td><td><pre>OnStart</pre></td><td>组件开始处理前</td><td>CallbackInput</td></tr>
|
||||
<tr><td><pre>TimingOnEnd</pre></td><td><pre>OnEnd</pre></td><td>组件成功返回后</td><td>CallbackOutput</td></tr>
|
||||
<tr><td><pre>TimingOnError</pre></td><td><pre>OnError</pre></td><td>组件返回错误时</td><td>error</td></tr>
|
||||
<tr><td><pre>TimingOnStartWithStreamInput</pre></td><td><pre>OnStartWithStreamInput</pre></td><td>组件接收流式输入时</td><td>StreamReader[CallbackInput]</td></tr>
|
||||
<tr><td><pre>TimingOnEndWithStreamOutput</pre></td><td><pre>OnEndWithStreamOutput</pre></td><td>组件返回流式输出时</td><td>StreamReader[CallbackOutput]</td></tr>
|
||||
<tr><td>TimingOnStart</td><td>OnStart</td><td>组件开始处理前</td><td>CallbackInput</td></tr>
|
||||
<tr><td>TimingOnEnd</td><td>OnEnd</td><td>组件成功返回后</td><td>CallbackOutput</td></tr>
|
||||
<tr><td>TimingOnError</td><td>OnError</td><td>组件返回错误时</td><td>error</td></tr>
|
||||
<tr><td>TimingOnStartWithStreamInput</td><td>OnStartWithStreamInput</td><td>组件接收流式输入时</td><td>StreamReader[CallbackInput]</td></tr>
|
||||
<tr><td>TimingOnEndWithStreamOutput</td><td>OnEndWithStreamOutput</td><td>组件返回流式输出时</td><td>StreamReader[CallbackOutput]</td></tr>
|
||||
</table>
|
||||
|
||||
**示例:ChatModel 调用流程**
|
||||
**非流式调用时序:**
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ ChatModel.Generate(ctx, messages) │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ OnStart │ ← 输入: CallbackInput (messages)
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 模型处理 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ OnEnd │ ← 输出: CallbackOutput (response)
|
||||
└──────────────────────┘
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as 业务代码
|
||||
participant CM as ChatModel
|
||||
participant CB as Callback Handler
|
||||
|
||||
Client->>CM: Generate(ctx, messages)
|
||||
CM->>CB: OnStart(messages)
|
||||
Note over CB: "记录输入,启动计时"
|
||||
CM->>CM: 模型处理
|
||||
CM->>CB: OnEnd(response)
|
||||
Note over CB: "记录输出,计算耗时"
|
||||
CM-->>Client: response
|
||||
```
|
||||
|
||||
**示例:流式输出流程**
|
||||
**流式调用时序:**
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ ChatModel.Stream(ctx, messages) │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ OnStart │ ← 输入: CallbackInput (messages)
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 模型处理(流式) │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ OnEndWithStreamOutput │ ← 输出: StreamReader[CallbackOutput]
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 逐个 chunk 返回 │
|
||||
└──────────────────────┘
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as 业务代码
|
||||
participant CM as ChatModel
|
||||
participant CB as Callback Handler
|
||||
|
||||
Client->>CM: Stream(ctx, messages)
|
||||
CM->>CB: OnStart(messages)
|
||||
Note over CB: "记录输入,启动计时"
|
||||
CM->>CM: 模型处理(流式)
|
||||
CM->>CB: OnEndWithStreamOutput(reader)
|
||||
Note over CB: "返回 StreamReader,逐 chunk 消费"
|
||||
loop 逐块消费
|
||||
CB->>CB: reader.Read()
|
||||
CB->>CB: 处理 chunk
|
||||
end
|
||||
CM-->>Client: stream chunks
|
||||
```
|
||||
|
||||
**注意:**
|
||||
> [!WARNING] 流式错误处理
|
||||
> 流式错误(stream 中途出错)**不会触发 OnError**,而是在 StreamReader 中返回。消费时务必检查 `reader.Err()`。
|
||||
|
||||
- 流式错误(stream 中途出错)不会触发 OnError,而是在 StreamReader 中返回
|
||||
- 同一 Handler 的 OnStart→OnEnd 可通过 context 传递状态
|
||||
- 不同 Handler 之间没有执行顺序保证
|
||||
### TimingChecker 优化
|
||||
|
||||
## Callback 的实现
|
||||
如果你的 Handler 不需要某些时机(比如只关心错误),可以实现 `TimingChecker` 接口来跳过不必要的调用开销:
|
||||
|
||||
### 1. 实现自定义 Callback Handler
|
||||
```go
|
||||
func (h *MyHandler) TimingChecker(timing callbacks.Timing) bool {
|
||||
// 只启用 Error 检测,其余跳过
|
||||
return timing == callbacks.TimingOnError
|
||||
}
|
||||
```
|
||||
|
||||
完整实现 `Handler` 接口需要实现所有 5 个方法,较为繁琐。Eino 提供了 `callbacks.HandlerHelper` 帮助类来简化实现:
|
||||
## Callback 实战
|
||||
|
||||
### 实现自定义 Callback Handler
|
||||
|
||||
直接实现全部 5 个方法比较繁琐。Eino 提供了 `callbacks.HandlerHelper` 链式构建器,只需注册感兴趣的回调:
|
||||
|
||||
```go
|
||||
import "github.com/cloudwego/eino/callbacks"
|
||||
@@ -246,122 +244,80 @@ handler := callbacks.NewHandlerHelper().
|
||||
callbacks.AppendGlobalHandlers(handler)
|
||||
```
|
||||
|
||||
**注意**:`RunInfo` 可能为 `nil`(如顶层调用没有 RunInfo),使用前请检查。
|
||||
**注意**:`RunInfo` 可能为 `nil`(如顶层调用),使用前务必检查。
|
||||
|
||||
### 2. 集成 CozeLoop
|
||||
### 集成与注册
|
||||
|
||||
```go
|
||||
func setupCozeLoop(ctx context.Context) (*cozeloop.Client, error) {
|
||||
apiToken := os.Getenv("COZELOOP_API_TOKEN")
|
||||
workspaceID := os.Getenv("COZELOOP_WORKSPACE_ID")
|
||||
|
||||
if apiToken == "" || workspaceID == "" {
|
||||
return nil, nil // 未配置则跳过
|
||||
}
|
||||
|
||||
client, err := cozeloop.NewClient(
|
||||
cozeloop.WithAPIToken(apiToken),
|
||||
cozeloop.WithWorkspaceID(workspaceID),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 注册为全局 Callback
|
||||
callbacks.AppendGlobalHandlers(clc.NewLoopHandler(client))
|
||||
|
||||
return client, nil
|
||||
}
|
||||
```
|
||||
|
||||
### 3. 在 main 中使用
|
||||
完整代码见 [cmd/ch06/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch06/main.go)。核心流程如下:
|
||||
|
||||
```go
|
||||
func main() {
|
||||
ctx := context.Background()
|
||||
|
||||
// 设置 CozeLoop(可选)
|
||||
client, err := setupCozeLoop(ctx)
|
||||
if err != nil {
|
||||
log.Printf("cozeloop setup failed: %v", err)
|
||||
}
|
||||
if client != nil {
|
||||
// 1. 可选:注册自定义日志 Callback
|
||||
handler := callbacks.NewHandlerHelper().
|
||||
OnStart(func(ctx context.Context, info *callbacks.RunInfo, input callbacks.CallbackInput) context.Context {
|
||||
log.Printf("[trace] %s/%s start", info.Component, info.Name)
|
||||
return ctx
|
||||
}).
|
||||
OnEnd(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
|
||||
log.Printf("[trace] %s/%s end", info.Component, info.Name)
|
||||
return ctx
|
||||
}).
|
||||
Handler()
|
||||
callbacks.AppendGlobalHandlers(handler)
|
||||
|
||||
// 2. 可选:启用 CozeLoop 链路追踪
|
||||
apiToken := os.Getenv("COZELOOP_API_TOKEN")
|
||||
workspaceID := os.Getenv("COZELOOP_WORKSPACE_ID")
|
||||
if apiToken != "" && workspaceID != "" {
|
||||
client, _ := cozeloop.NewClient(
|
||||
cozeloop.WithAPIToken(apiToken),
|
||||
cozeloop.WithWorkspaceID(workspaceID),
|
||||
)
|
||||
defer func() {
|
||||
time.Sleep(5 * time.Second) // 等待数据上报
|
||||
time.Sleep(5 * time.Second) // 等待数据上报
|
||||
client.Close(ctx)
|
||||
}()
|
||||
callbacks.AppendGlobalHandlers(clc.NewLoopHandler(client))
|
||||
}
|
||||
|
||||
// 创建 Agent 并运行...
|
||||
// 3. 正常创建并运行 Agent...
|
||||
}
|
||||
```
|
||||
|
||||
**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch06/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch06/main.go)):
|
||||
## 可观测性的三大价值
|
||||
|
||||
```go
|
||||
// 设置 CozeLoop 追踪
|
||||
cozeloopApiToken := os.Getenv("COZELOOP_API_TOKEN")
|
||||
cozeloopWorkspaceID := os.Getenv("COZELOOP_WORKSPACE_ID")
|
||||
if cozeloopApiToken != "" && cozeloopWorkspaceID != "" {
|
||||
client, err := cozeloop.NewClient(
|
||||
cozeloop.WithAPIToken(cozeloopApiToken),
|
||||
cozeloop.WithWorkspaceID(cozeloopWorkspaceID),
|
||||
)
|
||||
if err != nil {
|
||||
log.Fatalf("cozeloop.NewClient failed: %v", err)
|
||||
}
|
||||
defer func() {
|
||||
time.Sleep(5 * time.Second)
|
||||
client.Close(ctx)
|
||||
}()
|
||||
callbacks.AppendGlobalHandlers(clc.NewLoopHandler(client))
|
||||
}
|
||||
通过 Callback 收集的数据,我们可以实现三个层面的可观测性:
|
||||
|
||||
```mermaid
|
||||
quadrantChart
|
||||
title Observability Dimensions
|
||||
x-axis Low Impact --> High Impact
|
||||
y-axis Low Cost --> High Value
|
||||
"错误追踪": [0.8, 0.9]
|
||||
"成本优化": [0.6, 0.7]
|
||||
"性能分析": [0.4, 0.5]
|
||||
"审计合规": [0.9, 0.8]
|
||||
```
|
||||
|
||||
## 可观测性的价值
|
||||
|
||||
### 1. 性能分析
|
||||
|
||||
通过 Callback 收集的数据,可以分析:
|
||||
|
||||
- 模型调用延迟分布
|
||||
- Tool 执行时间排行
|
||||
- Token 消耗趋势
|
||||
|
||||
### 2. 错误追踪
|
||||
|
||||
当 Agent 出现问题时:
|
||||
|
||||
- 查看完整的调用链路
|
||||
- 定位是哪个环节出错
|
||||
- 分析错误原因
|
||||
|
||||
### 3. 成本优化
|
||||
|
||||
通过 Token 消耗数据:
|
||||
|
||||
- 识别高消耗的对话
|
||||
- 优化 Prompt 减少 Token
|
||||
- 选择更经济的模型
|
||||
| 维度 | 关键指标 | 典型场景 |
|
||||
|------|----------|----------|
|
||||
| **错误追踪** | 错误类型、堆栈、出错节点 | Agent 响应异常时快速定位是模型侧还是 Tool 侧的问题 |
|
||||
| **成本优化** | Token 消耗、每轮对话花费 | 识别高消耗对话,优化 Prompt 或切换更经济的模型 |
|
||||
| **性能分析** | 延迟分布、耗时 Top N | 发现慢查询——某个 Tool 执行时间过长影响整体体验 |
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Callback**:Eino 的观测钩子,在关键节点触发回调
|
||||
- **CozeLoop**:字节跳动的 AI 应用可观测性平台
|
||||
- **全局注册**:通过 `callbacks.AppendGlobalHandlers` 注册全局 Callback
|
||||
- **非侵入式**:业务代码不需要修改,Callback 自动触发
|
||||
- **可观测性价值**:性能分析、错误追踪、成本优化
|
||||
> [!SUMMARY] 要点回顾
|
||||
> - **Callback** 是 Eino 的非侵入式观测钩子,在组件生命周期的 5 个时机触发
|
||||
> - 使用 `callbacks.HandlerHelper` 可链式构建 Handler,只注册感兴趣的回调
|
||||
> - 通过 `callbacks.AppendGlobalHandlers` 注册全局 Callback,业务代码零修改
|
||||
> - **CozeLoop** 提供开箱即用的链路追踪和可视化
|
||||
> - 结合 `TimingChecker` 可实现性能最优的按需检测
|
||||
|
||||
## 扩展思考
|
||||
## 关联笔记
|
||||
|
||||
**其他 Callback 实现:**
|
||||
|
||||
- OpenTelemetry Callback:对接标准可观测性协议
|
||||
- 自定义日志 Callback:记录到本地文件
|
||||
- 指标 Callback:对接 Prometheus 等监控系统
|
||||
|
||||
**高级用法:**
|
||||
|
||||
- 在 Callback 中实现采样(只记录部分请求)
|
||||
- 在 Callback 中实现限流(根据 Token 消耗)
|
||||
- 在 Callback 中实现告警(错误率过高时通知)
|
||||
- [[Eino/quick_start/chapter_01_hello_eino.md]]
|
||||
- [[Eino/quick_start/chapter_04_tool.md]]
|
||||
- [[Eino/quick_start/chapter_05_agent.md]]
|
||||
|
||||
Reference in New Issue
Block a user