diff --git a/Eino/quick_start/chapter_01_chatmodel_and_message.md b/Eino/quick_start/chapter_01_chatmodel_and_message.md index 4dcb14e..442421e 100644 --- a/Eino/quick_start/chapter_01_chatmodel_and_message.md +++ b/Eino/quick_start/chapter_01_chatmodel_and_message.md @@ -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 { + <> + } + class ChatModel { + <> + +Generate() + +Stream() + } + class Tool { + <> + +Execute() + } + class Retriever { + <> + +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]] diff --git a/Eino/quick_start/chapter_02_chatmodelagent_runner_agentevent.md b/Eino/quick_start/chapter_02_chatmodelagent_runner_agentevent.md index e7360cf..fd74870 100644 --- a/Eino/quick_start/chapter_02_chatmodelagent_runner_agentevent.md +++ b/Eino/quick_start/chapter_02_chatmodelagent_runner_agentevent.md @@ -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` 如何支持流式消费。 + +--- + + + ## 代码位置 @@ -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["统一抽象
运行时多态"] :::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 完成所有事情? + - +
维度ChatModelChatModelAgent
定位Component(组件)Agent(智能体)
接口
Generate() / Stream()
Run() -> AsyncIterator[*AgentEvent]
输出直接返回消息内容返回事件流(包含消息、控制动作等)
输出直接返回消息内容返回事件流(含消息、控制动作等)
能力单纯的模型调用可扩展 tools、middleware、interrupt 等
适用场景简单的对话场景复杂的智能体应用
@@ -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["退出循环"] + + style S fill:#e1f5fe + style OUT fill:#ffebee +``` -**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [cmd/ch02/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch02/main.go)): +**逐步拆解:** + +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]] diff --git a/Eino/quick_start/chapter_03_memory_and_session.md b/Eino/quick_start/chapter_03_memory_and_session.md index 0eba844..ef891e1 100644 --- a/Eino/quick_start/chapter_03_memory_and_session.md +++ b/Eino/quick_start/chapter_03_memory_and_session.md @@ -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 ``` +--- + + + + ## 从内存到持久化:为什么需要 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()
保存用户消息"] + B --> C["session.GetMessages()
获取完整历史"] + C --> D["runner.Run(history)
Agent 处理消息"] + D --> E["收集助手回复"] + E --> F["session.Append()
保存助手消息"] -``` -┌─────────────────────────────────────────────────────────────┐ -│ 业务层(你的代码) │ -│ ┌─────────────┐ ┌──────────────┐ ┌───────────────┐ │ -│ │ 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
持久化存储"] + S2["GetMessages()"] + S3["Append()
保存消息"] + 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 文件存储方案适合简单的单机应用。在实际业务中,你可能需要考虑其他存储方案: -**其他存储实现:** + + + + + + + +
存储方案适用场景优势劣势
JSONL 文件单机应用、开发调试零依赖,简单直观不支持并发、分布式
SQLite / LevelDB桌面端应用轻量级嵌入式数据库不适合高并发写入
MySQL / PostgreSQL服务端部署成熟稳定,功能丰富运维成本较高
Redis分布式、高频访问性能极高,支持过期策略数据需额外持久化
S3 / OSS海量冷数据归档成本极低,无限扩展不适合频繁查询
-- 数据库存储(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]] diff --git a/Eino/quick_start/chapter_04_tool_and_filesystem.md b/Eino/quick_start/chapter_04_tool_and_filesystem.md index 9f0e389..b98a9ea 100644 --- a/Eino/quick_start/chapter_04_tool_and_filesystem.md +++ b/Eino/quick_start/chapter_04_tool_and_filesystem.md @@ -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) - - // 读取文件内容,支持按行偏移和限制 - Read(ctx context.Context, req *ReadRequest) (*FileContent, error) - - // 在文件中搜索匹配的内容 - GrepRaw(ctx context.Context, req *GrepRequest) ([]GrepMatch, error) - - // 根据 glob 模式匹配文件 - GlobInfo(ctx context.Context, req *GlobInfoRequest) ([]FileInfo, error) - - // 写入文件内容 - Write(ctx context.Context, req *WriteRequest) error - - // 编辑文件内容(字符串替换) - Edit(ctx context.Context, req *EditRequest) error + // 列出目录下的文件信息 + LsInfo(ctx context.Context, req *LsInfoRequest) ([]FileInfo, error) + + // 读取文件内容,支持按行偏移和限制 + Read(ctx context.Context, req *ReadRequest) (*FileContent, error) + + // 在文件中搜索匹配的内容 + GrepRaw(ctx context.Context, req *GrepRequest) ([]GrepMatch, error) + + // 根据 glob 模式匹配文件 + GlobInfo(ctx context.Context, req *GlobInfoRequest) ([]FileInfo, error) + + // 写入文件内容 + Write(ctx context.Context, req *WriteRequest) 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 对比:** - + - +
能力ChatModelAgentDeepAgent
能力ChatModelAgentDeepAgent
多轮对话✅✅
添加自定义 Tool✅ 手动注册每个 Tool✅ 手动注册或自动注册
文件系统访问(Backend)❌ 需手动创建并注册所有文件工具✅ 一级配置,自动注册
命令执行(StreamingShell)❌ 需手动创建✅ 一级配置,自动注册
内置任务管理❌✅
write_todos
工具
内置任务管理❌✅ `write_todos` 工具
支持子 Agent❌✅
-**选择建议:** - -- 纯对话场景(无外部访问)→ 用 `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 调用中间件与拦截器(下一章) diff --git a/Eino/quick_start/chapter_05_middleware.md b/Eino/quick_start/chapter_05_middleware.md index 0bae082..f4ae03d 100644 --- a/Eino/quick_start/chapter_05_middleware.md +++ b/Eino/quick_start/chapter_05_middleware.md @@ -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 生态还提供了以下开箱即用的中间件: + + + + + + +
Middleware功能说明
reduction工具输出缩减——当工具返回过长时自动截断并存入文件系统,防止上下文溢出
summarization对话历史摘要——Token 超阈值时自动生成摘要压缩历史,节省上下文空间
skill技能加载——让 Agent 按需动态加载预定义的 SKILL.md 知识包
+ +### 多 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:** - - - - - - -
Middleware功能说明
reduction工具输出缩减,当工具返回内容过长时自动截断并卸载到文件系统,防止上下文溢出
summarization对话历史自动摘要,当 token 数量超过阈值时自动生成摘要压缩历史
skill技能加载中间件,让 Agent 能够动态加载和执行预定义的技能
- -**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 可观测性(下一章) diff --git a/Eino/quick_start/chapter_06_callback_and_trace.md b/Eino/quick_start/chapter_06_callback_and_trace.md index 0ff7dfa..c6fdc9e 100644 --- a/Eino/quick_start/chapter_06_callback_and_trace.md +++ b/Eino/quick_start/chapter_06_callback_and_trace.md @@ -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*` - - - - - + + + + +
时机常量对应 Handler 方法触发点输入/输出
TimingOnStart
OnStart
组件开始处理前CallbackInput
TimingOnEnd
OnEnd
组件成功返回后CallbackOutput
TimingOnError
OnError
组件返回错误时error
TimingOnStartWithStreamInput
OnStartWithStreamInput
组件接收流式输入时StreamReader[CallbackInput]
TimingOnEndWithStreamOutput
OnEndWithStreamOutput
组件返回流式输出时StreamReader[CallbackOutput]
TimingOnStartOnStart组件开始处理前CallbackInput
TimingOnEndOnEnd组件成功返回后CallbackOutput
TimingOnErrorOnError组件返回错误时error
TimingOnStartWithStreamInputOnStartWithStreamInput组件接收流式输入时StreamReader[CallbackInput]
TimingOnEndWithStreamOutputOnEndWithStreamOutput组件返回流式输出时StreamReader[CallbackOutput]
-**示例: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]]