Eino学习笔记 - MuxiaoWF跳到主要内容

Eino学习笔记

Eino学习笔记 —— 入门篇

周一 7月 20 2026
7722 字 · 46 分钟

本笔记通过ai辅助生成,但对于代码阅读顺序应当没有问题(均使用eino官方的例子进行,个人认为官方的教程感觉有点太难懂了)

另附,可能是更新速度快,部分代码的方法已经被标记为废弃,但在示例中仍然未修改为最新的方法,建议自行对应

Eino 框架学习笔记

📂 源码仓库github.com/cloudwego/eino-examples

📖 系列文档入门笔记 | Agentic 进阶 | 附录一 Flow | 附录二 组件 | 附录三 工具


速览:两条调用路径 & 核心速记

在深入细节之前,先在心里放一张总地图。Eino 只有两条调用路径:

sequenceDiagram
    actor U as 用户代码
    participant CM as ChatModel
    participant A as Agent
    participant R as Runner
    participant I as Iterator

    Note over U,I: === 轻量路径:直接用 ChatModel(入门笔记 §2) ===
    U->>CM: Generate(messages) / Stream(messages)
    CM-->>U: *schema.Message / StreamReader

    Note over U,I: === 完整路径:Agent + Runner(入门笔记 §3~§5) ===
    U->>CM: ① NewChatModel(config)
    CM-->>U: model
    U->>A: ② NewChatModelAgent(model, instruction, tools)
    A-->>U: agent
    U->>R: ③ NewRunner(agent, streaming, checkpoint)
    R-->>U: runner
    U->>R: ④ Run(input) / Query(input, checkpointID)
    R-->>I: AsyncIterator[*AgentEvent]
    loop 消费事件
        I->>I: Next() → (event, ok)
    end

    Note over U,I: === Agentic 路径:Responses API(Agentic 进阶) ===
    U->>CM: ① New(agenticModel)
    CM-->>U: model
    U->>A: ② NewTypedChatModelAgent[*AgenticMessage](model, ...)
    A-->>U: agent
    U->>R: ③ NewTypedRunner[*AgenticMessage](agent, ...)
    R-->>U: runner
    U->>R: ④ Run(AgenticMessage输入)
    R-->>I: Iterator
    loop 消费事件
        I->>I: TypedGetMessage(event) → msg.String()
    end

📇 全篇速记卡片(看完这篇你只需要记住这些)

创建链:NewChatModel → NewChatModelAgent → NewRunner → Run
(模) (代) (跑) (消息)
两种迭代:Agent 用 Next() → (event, ok) —— !ok 结束
Stream 用 Recv() → (chunk, err) —— io.EOF 结束
Agent 管能力:{Name, Instruction, ToolsConfig, Model}
Runner 管执行:{Agent, EnableStreaming, CheckPointStore}
三种 Tool 创建:90% 用 InferTool(函数+struct tag)
需要约束用 NewTool(schema+函数)
有状态用结构体实现接口
中断恢复三要素:CheckPointStore + CheckPointID + Resume
核心包(按使用频率):
github.com/cloudwego/eino/adk ← Agent, Runner, Message
github.com/cloudwego/eino/schema ← UserMessage, SystemMessage
github.com/cloudwego/eino/components/tool/utils ← InferTool, NewTool
github.com/cloudwego/eino-ext/components/model/openai ← ChatModel

⚡ 关键概念:Iterator(事件迭代器)是什么?

在所有 Eino 示例中,runner.Run()runner.Query() 都会立即返回一个 AsyncIterator[*AgentEvent](简称 Iterator),然后你通过 Next() 逐个消费事件。这是 Eino 最核心的数据消费模式。

为什么要有 Iterator?

Agent 执行是异步的:模型调用可能耗时几秒到几十秒,期间可能经历”思考→调Tool→再思考→输出文本”多个阶段。如果 Run() 等全部完成才返回,你的程序就卡住了。

Iterator 解决了这个问题——立即返回一个”事件管道”,Agent 在后台生成事件、推入管道,你在前台逐个取出处理。类比:

Runner.Run() 返回 Iterator = 你拿到一个"对讲机"
Agent 在后台干活,每隔一会通过对讲机说句话
你在前台: Next() → 收到一条消息 → Next() → 收到下一条 → ...
Next() 返回 (event, false) = 对讲机没信号了 = Agent 干完了

Iterator 的本质

type AsyncIterator[T any] struct {
// 内部有 channel,连接"生产者"(Agent goroutine)和"消费者"(你的代码)
}
func (iter *AsyncIterator[T]) Next() (item T, ok bool) {
// 阻塞等待下一个事件
// ok=true: 拿到了新事件
// ok=false: 管道关闭,Agent 执行完毕
}

与 Go channel 的对比

Go channelAsyncIterator
创建make(chan T)Runner.Run() / Runner.Query() 返回
发送ch <- itemAgent 内部通过 gen.Send(event)
接收item := <-chevent, ok := iter.Next()
关闭close(ch)Agent 内部通过 gen.Close()
关闭检测item, ok := <-chevent, ok := iter.Next() — 同样是 ok=false

🧠 一句话:Iterator 就是包装了 channel 的”事件流”——Agent 在后台往里面写事件,你在前台用 Next() 逐个读。读完(ok=false)Agent 就结束了。


一、Hello World:认识四步组装线

📁 代码:adk/helloworld/helloworld.go

为什么这样设计?

Eino 把一个 AI 对话应用拆成四个独立零件,像组装流水线一样逐个创建、最后拼起来运行。这样拆的好处是:每个零件可以单独替换(比如换个模型、换个 Agent 类型),互不影响。

四步组装线(记忆口诀:“模→代→跑→消息”

第1步 创建 ChatModel → 第2步 创建 Agent → 第3步 创建 Runner → 第4步 发送消息并运行
("嘴巴",负责调用LLM) ("大脑",封装业务逻辑) ("引擎",管理执行过程) (用户说的话)

示范代码(带注释)

package main
import (
"context"
"fmt"
"log"
"os"
"github.com/cloudwego/eino/adk"
"github.com/cloudwego/eino/schema"
"github.com/cloudwego/eino-ext/components/model/openai"
)
func main() {
// 【固定写法】所有 Eino 调用的"通行证",贯穿整个调用链
ctx := context.Background()
// ─── 第1步:创建 ChatModel("嘴巴")───
// ChatModel = 与 LLM 通信的组件,屏蔽不同厂商差异
model, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{
APIKey: os.Getenv("OPENAI_API_KEY"),
Model: os.Getenv("OPENAI_MODEL"),
BaseURL: os.Getenv("OPENAI_BASE_URL"),
})
if err != nil {
log.Fatal(err)
}
// ─── 第2步:创建 Agent("大脑")───
agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "hello_agent",
Description: "A friendly greeting assistant",
Instruction: "You are a friendly assistant. Please respond warmly.",
Model: model, // ← 把第1步的 model 注入进来
})
if err != nil {
log.Fatal(err)
}
// ─── 第3步:创建 Runner("引擎")───
runner := adk.NewRunner(ctx, adk.RunnerConfig{
Agent: agent,
EnableStreaming: true,
})
// ─── 第4步:组装消息并运行 ───
input := []adk.Message{
schema.UserMessage("Hello, please introduce yourself."),
}
// runner.Run() 返回事件迭代器,逐个消费事件
events := runner.Run(ctx, input)
for {
event, ok := events.Next()
if !ok {
break
}
if event.Err != nil {
log.Printf("错误: %v", event.Err)
break
}
if msg, err := event.Output.MessageOutput.GetMessage(); err == nil {
fmt.Printf("Agent: %s\n", msg.Content)
}
}
}

💡 关键理解

概念一句话类比
context.Background()所有调用的”通行证”,固定写法身份证,每次办事都要出示
ChatModel负责与 LLM 通信的组件数据库驱动(屏蔽 MySQL/PG 差异)
Agent封装”怎么用模型”的业务逻辑业务逻辑层
Runner管理 Agent 的执行生命周期汽车引擎(你只转方向盘,不管内部)
events.Next()消费下一个事件,没事件了返回 false自助餐流水线,一盘盘拿

🧠 记忆技巧

把创建过程记成一句话:“用某个模型(ChatModel),造一个智能体(Agent),交给引擎(Runner)去跑(Run)”

NewChatModel → NewChatModelAgent → NewRunner → Run
(模) (代) (跑) (消息)

常用 import 路径:

"github.com/cloudwego/eino-ext/components/model/openai" // OpenAI 兼容模型,更多代码支持不再列出,直接参考官方文档
"github.com/cloudwego/eino/adk" // Agent, Runner, Message
"github.com/cloudwego/eino/schema" // UserMessage, SystemMessage 等

💡 类型别名说明adk.Message 就是 *schema.Message(Go type alias)。代码里写 []adk.Message{ schema.UserMessage(...) }[]*schema.Message{ schema.UserMessage(...) } 是等价的。文档里两种写法都会出现——见到时知道是同一个东西即可。

📇 本节速记卡片

NewChatModel → NewChatModelAgent → NewRunner → Run
event, ok := iter.Next(); if !ok { break }
msg, err := event.Output.MessageOutput.GetMessage()

二、直接用 ChatModel(不用 Agent)

📁 代码:quickstart/chat/main.gocomponents/prompt/chat_prompt/chat_prompt.go

什么时候不用 Agent?

如果只需要一次性的问答(不涉及多轮对话、工具调用、中断恢复),直接用 ChatModel 更轻量。

两种输出方式

方式方法效果适用场景
一次性model.Generate(ctx, messages)等全部生成完再返回 *schema.Message批量处理、非实时场景
流式model.Stream(ctx, messages)返回 *schema.StreamReader,逐块接收聊天、实时显示

示范代码

2.1 用模板构造消息

import (
"github.com/cloudwego/eino/components/prompt"
"github.com/cloudwego/eino/schema"
)
// 创建模板:用 {变量名} 做占位符
func createTemplate() prompt.ChatTemplate {
return prompt.FromMessages(schema.FString,
schema.SystemMessage("你是一个{role}。用{style}的语气回答。"),
schema.MessagesPlaceholder("chat_history", true),
schema.UserMessage("问题: {question}"),
)
}
// 用模板生成实际消息
func createMessages() []*schema.Message { // 正如上面所说,这里使用adk.Message也是一样的
template := createTemplate()
messages, err := template.Format(context.Background(), map[string]any{
"role": "程序员鼓励师",
"style": "积极、温暖且专业",
"question": "我的代码一直报错,该怎么办?",
"chat_history": []*schema.Message{
schema.UserMessage("你好"),
schema.AssistantMessage("嘿!加油!", nil),
},
})
// ... 处理 err
return messages
}

💡 components/prompt/chat_prompt/chat_prompt.go 展示了一个更完整的模板用法。

2.2 一次性生成(Generate)

result, err := model.Generate(ctx, messages)
if err != nil {
log.Fatal(err)
}
fmt.Println(result.Content)

2.3 流式生成(Stream)—— 重点掌握

streamReader, err := model.Stream(ctx, messages)
if err != nil {
log.Fatal(err)
}
defer streamReader.Close() // ⚠️ 用完记得关闭!
for {
chunk, err := streamReader.Recv()
if err == io.EOF { break }
if err != nil {
log.Fatal(err)
}
fmt.Print(chunk.Content)
}

💡 StreamReader 是什么? model.Stream() 返回的 *schema.StreamReader[*schema.Message] 和 Iterator 类似——模型在后台逐块生成,你通过 Recv() 逐块消费。但有关键区别:

  • Iterator(Agent 用):Next() 返回 (event, bool)!ok 表示 Agent 执行完毕
  • StreamReader(ChatModel 用):Recv() 返回 (chunk, error)io.EOF 表示流结束
  • defer Close() 是必须的:StreamReader 底层持有 HTTP 连接,不 Close 会导致连接泄漏 / goroutine 泄漏

⚠️ 本节最容易搞混的地方

Eino 里有两种迭代模式,结束判断条件不同:

// ─── 模式A:Agent Runner 事件迭代 ───
// Next() 返回 (event, bool),!ok 表示结束
for {
event, ok := iter.Next()
if !ok { break }
}
// ─── 模式B:ChatModel Stream 流式迭代 ───
// Recv() 返回 (chunk, error),io.EOF 表示结束
defer streamReader.Close()
for {
chunk, err := streamReader.Recv()
if err == io.EOF { break }
}

记忆诀窍:Agent 用 Next + ok,Stream 用 Recv + EOF。

常用 import 路径(本节新增):

"github.com/cloudwego/eino/components/prompt" // ChatTemplate, FromMessages

📇 本节速记卡片

一次性:model.Generate(ctx, messages) → (*Message, error)
流式: model.Stream(ctx, messages) → (*StreamReader, error)
defer close → Recv() 循环 → io.EOF 结束
模板: prompt.FromMessages(schema.FString, ...messages)
template.Format(ctx, map[string]any{...})

三、ChatModelAgent 基础:Tool + Agent + Runner

📁 代码:adk/intro/chatmodel/chatmodel.go

📖 理论:quickstart/chatwitheino/docs/ch02_chatmodel_agent_runner_console

🆕 本节新包速览(先记名字,后面会细讲):

"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool
"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool
"github.com/cloudwego/eino/compose" // ToolsNodeConfig(Tool 配置容器)

compose 包会在第七节深入讲解——这里你只需要知道 compose.ToolsNodeConfig 是注册 Tool 时用的”配置壳”即可。

3.1 为什么需要 Agent 而不是直接用 ChatModel?

维度ChatModel(组件)ChatModelAgent(智能体)
定位单个能力单元完整的 AI 应用
输出Generate() / Stream() 直接返回消息Run() 返回事件流,可包含工具调用、中断等
多轮对话需要自己管理 historyAgent 内部处理
工具调用❌ 不支持✅ 通过 ToolsConfig 配置
中断恢复❌ 不支持✅ 通过 CheckPointStore
适用场景简单的一次性问答复杂智能体应用

3.2 创建 Tool 的三种方式

📁 补充代码:quickstart/todoagent/main.go

Eino 提供三种创建 Tool 的方式:

方式做法适合场景
方式一:InferTool写一个普通函数,自动从 struct tag 推断 schema最常用,90% 场景
方式二:NewTool手动写 schema.ToolInfo,传给 utils.NewTool需要精确控制参数约束
方式三:结构体实现接口struct 实现 Info() + InvokableRun()Tool 逻辑复杂,需要带状态

方式一:InferTool(自动推断,最常用)⭐

import (
"github.com/cloudwego/eino/components/tool"
"github.com/cloudwego/eino/components/tool/utils"
)
// ① 定义输入结构体(用 json tag + jsonschema tag 描述参数)
type BookSearchInput struct {
Genre string `json:"genre" jsonschema_description:"书籍类型"`
MaxPages int `json:"max_pages" jsonschema_description:"最大页数"`
}
type BookSearchOutput struct {
Books []string `json:"books"`
}
// ② 写业务函数 + utils.InferTool 自动包装
func NewBookSearchTool() tool.InvokableTool {
t, _ := utils.InferTool(
"search_book",
"Search books by preferences",
func(ctx context.Context, input *BookSearchInput) (*BookSearchOutput, error) {
return &BookSearchOutput{Books: []string{"《三体》"}}, nil
},
)
return t
}

方式二:NewTool(手动定义 schema)

当需要更精确控制参数描述、添加枚举约束等场景时使用:

func getAddTodoTool() tool.InvokableTool {
info := &schema.ToolInfo{
Name: "add_todo",
Desc: "Add a todo item",
ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{
"content": {
Desc: "The content of the todo item",
Type: schema.String,
Required: true,
},
"deadline": {
Desc: "The deadline of the todo item, in unix timestamp",
Type: schema.Integer,
},
}),
}
return utils.NewTool(info, AddTodoFunc)
}
type TodoAddParams struct {
Content string `json:"content"`
Deadline *int64 `json:"deadline,omitempty"`
}
func AddTodoFunc(_ context.Context, params *TodoAddParams) (string, error) {
return `{"msg": "add todo success"}`, nil
}

方式三:结构体实现接口(最灵活)

Tool 需要带状态、或者逻辑很复杂时,直接实现接口:

type ListTodoTool struct{}
func (lt *ListTodoTool) Info(_ context.Context) (*schema.ToolInfo, error) {
return &schema.ToolInfo{
Name: "list_todo",
Desc: "List all todo items",
ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{
"finished": {
Desc: "filter todo items if finished",
Type: schema.Boolean,
},
}),
}, nil
}
func (lt *ListTodoTool) InvokableRun(_ context.Context, argumentsInJSON string, _ ...tool.Option) (string, error) {
return `{"todos": [...]}`, nil
}

🧠 怎么选:90% 场景用 InferTool 就够了;需要手动约束参数类型/枚举值时用 NewTool;Tool 本身带状态(比如有数据库连接)时用结构体方式。三种方式注册到 Agent 的方式完全一样——都丢进 ToolsConfig.Tools 数组里。

3.3 组装完整应用(Tool + 流式)

ℹ️ 本节聚焦:只展示 Tool + Agent + Runner + Query 的主流程。中断恢复(Resume)会单独在第五节讲。

package main
import (
"context"
"fmt"
"log"
"os"
"github.com/cloudwego/eino/adk"
"github.com/cloudwego/eino/components/tool"
"github.com/cloudwego/eino/compose"
"github.com/cloudwego/eino/schema"
"github.com/cloudwego/eino-ext/components/model/openai"
)
func main() {
ctx := context.Background()
// ─── 创建模型 ───
model, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{
APIKey: os.Getenv("OPENAI_API_KEY"),
Model: os.Getenv("OPENAI_MODEL"),
BaseURL: os.Getenv("OPENAI_BASE_URL"),
})
if err != nil {
log.Fatal(err)
}
// ─── 创建 Agent(配置 Tool)───
agent, err := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "BookRecommender",
Description: "An agent that recommends books",
Instruction: "You are a book expert. Use search_book tool to find books.",
Model: model,
ToolsConfig: adk.ToolsConfig{
ToolsNodeConfig: compose.ToolsNodeConfig{
Tools: []tool.BaseTool{
NewBookSearchTool(),
},
},
},
})
if err != nil {
log.Fatal(err)
}
// ─── 创建 Runner(配置流式)───
runner := adk.NewRunner(ctx, adk.RunnerConfig{
Agent: agent,
EnableStreaming: true,
})
// ─── 运行 ───
iter := runner.Query(ctx, "推荐一本科幻小说")
for {
event, ok := iter.Next()
if !ok { break }
if event.Err != nil {
log.Fatal(event.Err)
}
fmt.Printf("Event: %+v\n", event)
}
}

3.4 配置项归属速查

记住一句话:Agent 管”有什么能力”,Runner 管”怎么执行”。

配置在哪里设置作用
InstructionChatModelAgentConfig系统提示词,定义 Agent 的”人设”
ToolsConfigChatModelAgentConfig注册 Agent 可用的工具
EnableStreaming: trueRunnerConfig流式输出(打字机效果)
CheckPointStoreRunnerConfig保存执行状态,支持中断恢复(第五节)

3.5 底层视角:Agent 的本质是一个接口

📁 代码:adk/intro/custom/myagent.go

adk.NewChatModelAgent 内部实现了 adk.Agent 接口:

type Agent interface {
Name(ctx) string
Description(ctx) string
Run(ctx, input, ...options) *AsyncIterator[*AgentEvent]
}

myagent.go 展示了最简实现——三个关键点:

  • NewAsyncIteratorPair = 建一条事件流水线,gen 管发送,iter 管接收
  • 开 goroutine 是因为 Run() 要立即返回 iter,事件在后台异步生成
  • gen.Close() 必须调用,否则消费端的 Next() 会永远阻塞

常用 import 路径(本节新增):

"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool, Option
"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool
"github.com/cloudwego/eino/compose" // ToolsNodeConfig

📇 本节速记卡片

Agent.{Name, Instruction, ToolsConfig, Model}
Runner.{Agent, EnableStreaming, CheckPointStore}
Tool 三种创建:InferTool(90%)/ NewTool / 结构体接口
注册方式都一样:ToolsConfig.Tools = []tool.BaseTool{ toolfuction...}
附注:如果忘记某特定的struct需要填写什么可以直接看源码接口定义

四、Middleware:给 Agent 加拦截器

📁 代码:

  • quickstart/chatwitheino/helpers/middleware.go - SafeToolMiddleware,自定义一个新的middleware
  • quickstart/chatwitheino/helpers/retry.go - 错误自动重试
  • adk/middlewares/dynamictool/toolsearch/ — 动态工具检索中间件
  • adk/middlewares/skill/main.go — 技能加载中间件

📖 理论:quickstart/chatwitheino/docs/ch05_middleware

4.1 为什么需要 Middleware?

两个常见问题:

问题一:Tool 报错导致整个对话崩溃

[tool call] read_file(file_path: "nonexistent.txt")
Error: open nonexistent.txt: no such file or directory
// 💥 对话直接中断

问题二:模型 API 限流导致失败

Error: rate limit exceeded (429)
// 💥 对话中断

你希望的行为是:Tool 报错时把错误交给 LLM 自愈模型限流时自动重试。这就是 Middleware 要解决的问题——Agent 的拦截器,在调用前后插入自定义逻辑。

4.2 Middleware 的本质:洋葱模型

请求 → A.Wrap → B.Wrap → C.Wrap → 实际执行 → C返回 → B返回 → A返回 → 响应
↑ ↑
最外层先拦截 最内层先接触实际执行结果

4.3 核心场景一:SafeToolMiddleware(自定义一个新的middleware:Tool 错误转字符串)

// 定义一个结构体来承载中间件的逻辑。内嵌一个基础中间件(Base Middleware),以复用其默认行为。
type safeToolMiddleware struct {
*adk.BaseChatModelAgentMiddleware
}
// 拦截同步工具调用,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) {
// 拦截原始 endpoint的结果和err
result, err := endpoint(ctx, args, opts...)
if err != nil {
// ⚠️ 中断错误必须继续传播,不能吞掉
if _, ok := compose.IsInterruptRerunError(err); ok {
return "", err
}
// 普通错误 → 转成字符串,让 LLM 看到
return fmt.Sprintf("[tool error] %v", err), nil
}
// 没有错误就直接返回原始的结果
return result, nil
}, nil
}
// 还可以拦截其他的调用如WrapStreamableToolCall 拦截流式工具调用等等,可以自行查看官方示例

💡 关键区分compose.IsInterruptRerunError 是中断恢复(第五节)抛出的特殊错误,它必须继续向上传播,不能转成字符串。

4.4 核心场景二:ModelRetryConfig(模型调用自动重试)

原来的IsRetryAble被标记弃用,现在应该使用ShouldRetry进行替代

部分中间件可以直接在agent定义时注明,如ModelRetryConfig

ModelRetryConfig: &adk.ModelRetryConfig{
MaxRetries: 5,
ShouldRetry: func(ctx context.Context, retryCtx *adk.TypedRetryContext[M]) *adk.TypedRetryDecision[M] {
// 获取错误信息
err := retryCtx.Error
// 判断是否需要重试
if err != nil && (strings.Contains(err.Error(), "429") || strings.Contains(err.Error(), "Too Many Requests")) {
// 返回重试决策
return &adk.TypedRetryDecision[M]{
Retry: true,
// 在这里你还可以修改重试时的输入或选项,等等
}
}
// 不需要重试,返回 nil 或 Retry: false
return nil
},
},

4.5 核心场景三:动态工具检索(ToolSearch Middleware)

eino还预先定义了部分middleware,如agentsmd、dynamictool/toolsearch等,可以直接调用,下面介绍一些常见的

📁 adk/middlewares/dynamictool/toolsearch/

痛点:Tool 太多会撑爆上下文窗口。ToolSearch 先把所有 Tool 放进”工具库”,调用时搜索过滤:

toolSearchMiddleware, _ := toolsearch.New(ctx, &toolsearch.Config{
DynamicTools: allDynamicTools,
})
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "tool_search_agent",
Model: chatModel,
Handlers: []adk.ChatModelAgentMiddleware{
toolSearchMiddleware, // ← 不是 ToolsConfig!
},
})

🧠 与第三节 ToolsConfig 的区别:第三节把 Tool 写在 ToolsConfig.Tools 里——全部发给 LLM。ToolSearch 先用搜索过滤,只把相关的发给 LLM。

4.6 核心场景四:Skill Middleware(动态加载技能)

📁 adk/middlewares/skill/main.go

创建 Skill 后端

skillBackend, err := skill.NewBackendFromFilesystem(ctx, &skill.BackendFromFilesystemConfig{
Backend: be,
BaseDir: skillsDir, // 指定扫描的根目录
})

创建并注入 Skill 中间件

sm, err := skill.NewMiddleware(ctx, &skill.Config{
Backend: skillBackend,
})

之后在运行时从文件系统动态加载”技能”(预定义的提示词+工具组合)

// 文件系统中间件 → Skill 中间件(有序!)
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "LogAnalysisAgent",
Model: cm,
Handlers: []adk.ChatModelAgentMiddleware{
fsm, // 文件读写能力(洋葱模型外层,最先wrap,最后返回)
sm, // Skill 加载能力(洋葱模型内层,之后wrap,先返回)
},
})

4.7 Middleware 全景速查

Middleware 类型解决什么痛点来源配置位置
SafeToolMiddlewareTool 报错中断对话 → 转成字符串让 LLM 自愈自己实现Handlers
ModelRetryConfig模型 API 限流 → 自动重试adk.ModelRetryConfig(内置)ChatModelAgentConfig
ToolSearchTool 太多撑爆上下文 → 动态检索注入toolsearch.New()Handlers
Skill运行时动态加载技能 → 从文件系统读skill.NewMiddleware()Handlers
filesystemAgent 需要读写文件系统filesystem.New()Handlers

🧠 模型:安检通道

用户请求 → [安检A:查证件] → [安检B:过X光] → [安检C:金属探测] → 登机
用户收到 ← [安检A:盖章] ← [安检B:贴标签] ← [安检C:称重] ← 行李出来

📇 本节速记卡片

SafeToolMiddleware = 错误转字符串(除了 InterruptRerunError)
ModelRetryConfig = 限流自动重试(内置配置项)
ToolSearch = 工具池动态搜索(替代 ToolsConfig)
Skill = 文件系统加载技能
通用模式:嵌入 BaseChatModelAgentMiddleware → 覆写一个钩子 → 放入 Handlers

五、中断与恢复(Interrupt & Resume)

📁 代码:

  • adk/intro/chatmodel/ — Tool 中断
  • adk/cancel/graceful-exit/ — 外部取消
  • adk/human-in-the-loop/ — 人机协作全套示例

📖 Memory 理论:quickstart/chatwitheino/docs/ch03_memory_session_jsonl

📖 Tool 理论:quickstart/chatwitheino/docs/ch04_tool_backend_filesystem

5.1 中断机制总览

这是 Eino ADK 的核心高级特性。整个机制基于一个简单的状态机:

stateDiagram-v2
    [*] --> Running: runner.Query/Run(checkpointID)
    Running --> Interrupted: compose.Interrupt()
    Running --> Cancelled: Ctrl-C → cancelFn()
    Running --> Completed: 正常结束
    Running --> HITL_Interrupted: compose.Interrupt()

    Interrupted --> Resuming: runner.Resume(checkpointID, WithToolOptions)
    Cancelled --> ResumingNew: 新runner.Resume(checkpointID)
    HITL_Interrupted --> ResumingHITL: runner.ResumeWithParams(checkpointID, params)

    Resuming --> Running
    ResumingNew --> Running
    ResumingHITL --> Running
    Completed --> [*]

三种触发方式对比如下:

方式一:Tool 中断(本质就是方式三)方式二:外部取消方式三:Human-in-the-Loop
触发源Tool 调用 `compose.Interrupt“系统信号 Ctrl-CTool 调用 compose.Interrupt
Runner同一个 runner新建 runner同一个 runner
Resume APIrunner.Resume + WithToolOptionsrunner.Resumerunner.ResumeWithParams
典型场景信息不足,等用户补充长时间运行,用户想中止敏感操作审批、参数审核

共同点:都依赖 CheckPointStore + CheckPointID

🆕 本节新包速览

"github.com/cloudwego/eino/compose" // CheckPointStore(中断状态存储接口)
// CheckPointStore 有多种实现,如:
// store.NewInMemoryStore() ← 内存实现(开发用)
// 第五节只关注中断API,CheckPointStore 的底层原理见第七节

5.2 方式一:Tool 内部中断(主动暂停等输入)

用户输入 "推荐一本书"
runner := adk.NewRunner(ctx, adk.RunnerConfig{
Agent: a,
CheckPointStore: store.NewInMemoryStore(), // runner设置CheckPointStore以储存后续的CheckPointID
})
runner.Query(ctx, "推荐一本书", WithCheckPointID("1")) 传入CheckPointID,后续恢复通过这个
Agent 分析:信息不够 → Tool 返回 Interrupt("请问你喜欢什么类型的书?") 注:示例代码中NewInterruptAndRerunErr被标记启用
事件流中收到 Interrupted 事件 → 暂停,等待用户输入
runner.Resume(ctx, "1", WithToolOptions(WithNewInput("科幻"))) 通过CheckPointID进行恢复
Agent 继续执行 → 用 "科幻" 重新调用 Tool → 返回结果

注:既然这个已经过时+本质和方法三相同,可以直接看方法三是如何中断的

5.3 方式二:外部取消(Ctrl-C 优雅退出)

📁 adk/cancel/graceful-exit/main.go

设计意图:用户按 Ctrl-C 时不直接杀进程白跑——而是安全暂停、自动保存进度、之后从断点恢复。

用户按 Ctrl-C
cancelFn(CancelAfterChatModel, WithRecursive, 30s超时)
├─→ CancelAfterChatModel: 等当前 ChatModel 调用完成再停
├─→ WithRecursive: 递归传播到子 Agent
└─→ 超时兜底: 30s 内没到安全点 → 强制 CancelImmediate
CheckPointStore 自动保存状态 → 新 Runner.Resume 恢复

取消的创建

// 创建取消选项与外部触发函数
cancelOpt, cancelFn := adk.WithCancel()
// 启动Agent运行,绑定断点ID
iter := runner.Run(ctx, input, cancelOpt, adk.WithCheckPointID(checkpointID))
// 注册系统信号监听:SIGINT(Ctrl-C)、SIGTERM
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
go func() {
sig := <-sigCh
fmt.Printf("\n收到操作系统信号: %v,发起Agent优雅取消\n", sig)
// 发起优雅取消:等待ChatModel安全点、递归传播子Agent、30s超时兜底
handle, contributed := cancelFn(
adk.WithAgentCancelMode(adk.CancelAfterChatModel),
adk.WithRecursive(),
adk.WithAgentCancelTimeout(30*time.Second),
)
fmt.Printf("取消请求是否成功提交: %v\n", contributed)
// 阻塞等待取消流程完成(断点持久化完成后返回)
if waitErr := handle.Wait(); waitErr != nil {
fmt.Printf("优雅取消异常: %v\n", waitErr)
} else {
fmt.Println("优雅取消完成,断点已保存")
}
}()

取消的恢复

// 使用同一个CheckPointStore,新建Runner恢复
resumeRunner := adk.NewRunner(ctx, adk.RunnerConfig{
Agent: agent,
EnableStreaming: true,
CheckPointStore: cpStore,
})
resumeIter, err := resumeRunner.Resume(ctx, checkpointID)
if err != nil {
log.Fatalf("断点恢复失败: %v", err)
}
drainEvents(resumeIter)

5.4 方式三:Human-in-the-Loop 人机协作(ResumeWithParams)

📁 adk/human-in-the-loop/1_approval/ ~ 8_supervisor-plan-execute/

通过 interruptID 精确匹配中断点,注入任意类型的数据:

// ① 启动
iter := runner.Query(ctx, query, adk.WithCheckPointID("1"))
// ② 检测中断,提取 interruptID
interruptCtx := lastEvent.Action.Interrupted.InterruptContexts[0]
interruptID := interruptCtx.ID
// ③ 用户决策后,通过 ResumeWithParams 传回
iter, _ = runner.ResumeWithParams(ctx, "1", &adk.ResumeParams{
Targets: map[string]any{
interruptID: modifiedInfo,
},
})

中断方:

func FollowUp(ctx context.Context, input *FollowUpToolInput) (string, error) {
// 首次进入(触发中断)
wasInterrupted, _, storedState := tool.GetInterruptState[*FollowUpState](ctx)
if !wasInterrupted {
// 准备要展示给用户的信息(比如问题列表)
info := &FollowUpInfo{Questions: input.Questions}
// 准备要保存的状态(为了恢复时知道我们在干嘛)
state := &FollowUpState{Questions: input.Questions}
// 触发中断!把 info 发给前端展示,把 state 存起来,然后函数立刻返回
return "", tool.StatefulInterrupt(ctx, info, state)
}
// 第二次进入
isResumeTarget, hasData, resumeData := tool.GetResumeContext[*FollowUpInfo](ctx)
if !isResumeTarget {
// 如果系统唤醒了这个函数,但发现它不是当前该恢复的目标
// 那就把之前保存的状态再拿出来,继续中断等待
info := &FollowUpInfo{Questions: storedState.Questions}
return "", tool.StatefulInterrupt(ctx, info, storedState)
}
// 被唤醒,但还没拿到答案(防御性重入)
if !hasData || resumeData.UserAnswer == "" {
return "", fmt.Errorf("tool resumed without a user answer")
}
// 成功恢复,拿到答案(正常退出)
return resumeData.UserAnswer, nil
}

四种交互模式(区别仅在于 info 类型和用户操作):

模式中断时携带的信息用户操作示例目录
审批ApprovalResultY/N 批准或拒绝1_approval, 5_supervisor
审核编辑ReviewEditInfo修改参数 / 批准 / 拒绝2_review-and-edit
反馈循环FeedbackInfo输入反馈意见3_feedback-loop
追问FollowUpInfo逐个回答问题4_follow-up, 7_deep-agents

💡 ResumeWithParams vs 旧 ResumeResumeWithParams 是通用 API,Resume + WithToolOptions 是其特例。

📇 本节速记卡片

中断恢复三要素:CheckPointStore(存状态)+ CheckPointID(标识)+ Resume(恢复)
方式一:compose.Interrupt → runner.Resume(id, WithToolOptions)
方式二:Ctrl-C → cancelFn → 新 runner.Resume(id)
方式三:compose.Interrupt → runner.ResumeWithParams(id, {Targets: {interruptID: data}})

六、进阶主题

6.1 Callbacks 与 Tracing(可观测性)

📁 adk/common/trace/coze_loop.go

traceCloseFn, startSpanFn := trace.AppendCozeLoopCallbackIfConfigured(ctx)
defer traceCloseFn(ctx)
ctx, endSpanFn := startSpanFn(ctx, "MyTask", "user query")
// ... runner.Query(ctx, query) ...
endSpanFn(ctx, lastMessage)

6.2 Workflow Agent:编排多个 Agent 的协作方式

📁 adk/intro/workflow/loop/, parallel/, sequential/

三种内置工作流 Agent:

类型模式适用场景
LoopAgent主Agent生成 → 评审审查 → 不满意重改需要反复打磨的任务
ParallelAgent多个子Agent同时接收相同输入,各自执行同时从多数据源收集信息
SequentialAgent子Agent逐个执行,前一个输出→后一个输入分步任务流水线
// LoopAgent
a, _ := adk.NewLoopAgent(ctx, &adk.LoopAgentConfig{
SubAgents: []adk.Agent{mainAgent, critiqueAgent},
MaxIterations: 5,
})
// ParallelAgent
a, _ := adk.NewParallelAgent(ctx, &adk.ParallelAgentConfig{
SubAgents: []adk.Agent{stockAgent, newsAgent, socialAgent},
})
// SequentialAgent
a, _ := adk.NewSequentialAgent(ctx, &adk.SequentialAgentConfig{
SubAgents: []adk.Agent{planAgent, writerAgent},
})

6.3 Multi-Agent:Supervisor 模式

📁 adk/multiagent/supervisor/, layered-supervisor/, plan-execute-replan/

Supervisor:一个”主管”Agent 分析用户请求,决定派给哪个”专家”子 Agent。

sv, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "supervisor",
Instruction: "You are a supervisor managing a research_agent and math_agent.",
Model: m,
})
supervisorAgent, _ := supervisor.New(ctx, &supervisor.Config{
Supervisor: sv,
SubAgents: []adk.Agent{searchAgent, mathAgent},
})

Layered Supervisor:子 Agent 本身也可以是 Supervisor,形成树状调度。

Plan-Execute-Replan:Planner 制定计划 → Executor 逐步执行 → Replanner 检查并修改计划。

entryAgent, _ := planexecute.New(ctx, &planexecute.Config{
Planner: planAgent,
Executor: executeAgent,
Replanner: replanAgent,
MaxIterations: 20,
})

🧠 三种多 Agent 模式对比

  • Workflow Agent(6.2):编排是固定的(顺序/并行/循环)
  • Supervisor(6.3):编排是动态的,由 LLM 决定派给谁
  • Plan-Execute-Replan:适合需要多步规划且计划可能中途调整的复杂任务

6.4 Compose 编排 API 总览

按抽象层级从高到低:

ADK (Agent+Runner) > Chain > Graph > Workflow
高层,开箱即用 底层,精细控制每个节点和数据流
API特点适合场景
Workflow字段映射 + 自动依赖推导纯函数数据流水线
Graph手动 AddEdge 建图,支持 StateLLM 调用流程、有状态分支
Chain链式 Append 节点ChatModel → ToolsNode 线性流程
Batch批量并行处理 N 个输入文档批量审核等

Graph 基础用法

g := compose.NewGraph[map[string]any, *schema.Message]()
_ = g.AddChatTemplateNode("prompt", pt)
_ = g.AddChatModelNode("model", chatModel)
_ = g.AddEdge(compose.START, "prompt")
_ = g.AddEdge("prompt", "model")
_ = g.AddEdge("model", compose.END)
runner, _ := g.Compile(ctx)
result, _ := runner.Invoke(ctx, input)

Chain 基础用法

chain := compose.NewChain[map[string]any, string]()
chain.
AppendLambda(preprocess).
AppendBranch(compose.NewChainBranch(cond).AddLambda("b1", fn1).AddLambda("b2", fn2)).
AppendLambda(postprocess)
runner, _ := chain.Compile(ctx)

📇 本节速记卡片

Workflow Agent: LoopAgent / ParallelAgent / SequentialAgent
Supervisor: supervisor.New({Supervisor, SubAgents})
Compose 编排: Graph(AddEdge) / Chain(Append) / Workflow(MapFields)

七、完整架构分层图

graph TD
    Q["你要做什么?"] --> Q1{"简单一次性问答?"}
    Q1 -->|是| CM["ChatModel.Generate/Stream
§2"] Q1 -->|否| Q2{"标准 Agent
问答+工具+记忆?"} Q2 -->|是| ADK["ChatModelAgent + Runner
§3"] Q2 -->|否| Q3{"自定义 ReAct 循环?"} Q3 -->|是| REACT["react.NewAgent
appendix1 §1"] Q3 -->|否| Q4{"多 Agent 协作?"} Q4 -->|调度模式| SUP["ADK Supervisor
§6.3
或 Host MultiAgent
appendix1 §6"] Q4 -->|流水线| PE["Plan-Execute
appendix1 §7"] Q4 -->|复杂状态图| CUSTOM["自定义 Graph+State
appendix1 §9"] Q4 -->|否| Q5{"精细控制每个节点?"} Q5 -->|是| GRAPH["Compose Graph/Chain
§6.4"] Q5 -->|否| CHECK["回到 ADK 路径"]

📇 全篇速记卡片(完整版)

┌─ 创建链 ─────────────────────────────────────────────┐
│ NewChatModel → NewChatModelAgent → NewRunner → Run │
│ (模) (代) (跑) (消息) │
├─ 两种迭代 ───────────────────────────────────────────┤
│ Agent: for { event, ok := iter.Next(); if !ok {break}}│
│ Stream: for { chunk, err := sr.Recv(); io.EOF → break}│
├─ 职责分离 ───────────────────────────────────────────┤
│ Agent 管能力:{Name, Instruction, ToolsConfig, Model} │
│ Runner管执行:{Agent, EnableStreaming, CheckPointStore}│
├─ Tool 创建 ──────────────────────────────────────────┤
│ InferTool(90%)/ NewTool / 结构体接口 │
├─ Middleware ────────────────────────────────────────┤
│ SafeTool / ModelRetry / ToolSearch / Skill │
├─ 中断恢复 ──────────────────────────────────────────┤
│ CheckPointStore + CheckPointID + Resume │
│ 三种方式:Tool中断 / Ctrl-C取消 / Human-in-the-Loop │
├─ 编排层次 ──────────────────────────────────────────┤
│ ADK(Agent+Runner) > Chain > Graph > Workflow │
├─ Agentic 路径([Agentic 进阶](./agentic))──────────────────────────┤
│ 所有类型加 Typed[*AgenticMessage] 即可 │
└──────────────────────────────────────────────────────┘

常用 import 路径(全篇汇总):

// 模型
"github.com/cloudwego/eino-ext/components/model/openai" // OpenAI 兼容 ChatModel
// ADK 核心
"github.com/cloudwego/eino/adk" // Agent, Runner, Message, Middleware
"github.com/cloudwego/eino/schema" // UserMessage, SystemMessage, ToolInfo
// Tool
"github.com/cloudwego/eino/components/tool" // BaseTool, InvokableTool
"github.com/cloudwego/eino/components/tool/utils" // InferTool, NewTool
// 编排
"github.com/cloudwego/eino/compose" // Graph, Chain, Workflow, ToolsNodeConfig
// Middleware 扩展
"github.com/cloudwego/eino/adk/middlewares/dynamictool/toolsearch" // ToolSearch
"github.com/cloudwego/eino/adk/middlewares/skill" // Skill
"github.com/cloudwego/eino/adk/middlewares/filesystem" // 文件系统
// Prompt
"github.com/cloudwego/eino/components/prompt" // ChatTemplate, FromMessages

⚠️ 常见坑(Common Gotchas)

开发 Eino 应用时,下面几个坑最容易踩到——每个都用一个”坏代码 vs 好代码”的对比来展示:

坑 1:忘记 defer streamReader.Close() → goroutine 泄漏

// ❌ 坏:忘记 Close,HTTP 连接永不释放,goroutine 泄漏
sr, _ := model.Stream(ctx, messages)
for { chunk, err := sr.Recv(); ... }
// ✅ 好:defer Close 保证无论怎么退出都会关闭
sr, _ := model.Stream(ctx, messages)
defer sr.Close()
for { chunk, err := sr.Recv(); ... }

坑 2:自定义 Agent 忘记 gen.Close()Next() 永远阻塞

// ❌ 坏:忘记 Close,消费端 Next() 永远等不到 false
func (m *MyAgent) Run(...) *AsyncIterator[*AgentEvent] {
iter, gen := adk.NewAsyncIteratorPair[*AgentEvent]()
go func() {
gen.Send(&AgentEvent{...})
// 忘记 gen.Close()!!!
}()
return iter
}
// ✅ 好:defer + recover 确保 Close 一定被调用
func (m *MyAgent) Run(...) *AsyncIterator[*AgentEvent] {
iter, gen := adk.NewAsyncIteratorPair[*AgentEvent]()
go func() {
defer func() {
if e := recover() { gen.Send(&AgentEvent{Err: ...}) }
gen.Close() // ← 必须调用!
}()
gen.Send(&AgentEvent{...})
}()
return iter
}

坑 3:SafeToolMiddleware 误吞中断错误

// ❌ 坏:所有 error 都转成字符串,中断错误也被吞了
func wrap(endpoint) endpoint {
return func(ctx, args) (string, error) {
result, err := endpoint(ctx, args)
if err != nil { return fmt.Sprintf("[error] %v", err), nil }
return result, nil
}
}
// ✅ 好:区分中断错误和普通错误
func wrap(endpoint) endpoint {
return func(ctx, args) (string, error) {
result, err := endpoint(ctx, args)
if err != nil {
if _, ok := compose.IsInterruptRerunError(err); ok {
return "", err // ← 中断错误必须传播!
}
return fmt.Sprintf("[tool error] %v", err), nil
}
return result, nil
}
}

坑 4:Resume 时 CheckPointID 不一致

// ❌ 坏:Query 和 Resume 用了不同的 ID——Runner 找不到之前的保存点
iter := runner.Query(ctx, "帮我推荐", adk.WithCheckPointID("session-1"))
// ... 中断 ...
iter, _ = runner.Resume(ctx, "session-2") // ← ID 不一致!应该是 "session-1"
// ✅ 好:同一个 ID
iter := runner.Query(ctx, "帮我推荐", adk.WithCheckPointID("session-1"))
// ... 中断 ...
iter, _ = runner.Resume(ctx, "session-1") // ← 一致!

坑 5:Resume 时用了新的 Runner 但没传同一个 CheckPointStore

// ❌ 坏:新 Runner 没用同一个 CheckPointStore——状态丢失
runner1 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent}) // 默认空 Store
iter := runner1.Query(ctx, "query", adk.WithCheckPointID("1"))
// ... Ctrl-C 取消 ...
runner2 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent}) // 另一个空 Store
iter, _ = runner2.Resume(ctx, "1") // ← 找不到 checkpoint!
// ✅ 好:共享同一个 CheckPointStore
store := store.NewInMemoryStore()
runner1 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent, CheckPointStore: store})
runner1.Query(...)
// ...
runner2 := adk.NewRunner(ctx, adk.RunnerConfig{Agent: agent, CheckPointStore: store})
runner2.Resume(ctx, "1") // ← 找得到!

坑 6:服务端工具(web_search)期待 function_tool_result

// ❌ 坏:等不到 web_search 的 function_tool_result
for _, block := range msg.ContentBlocks {
if block.Type == schema.ContentBlockTypeFunctionToolResult {
// web_search 永远不会产生这个!
}
}
// ✅ 好:理解 server_tool_call 在厂商侧执行,没有对应的 function_tool_result
for _, block := range msg.ContentBlocks {
if block.Type == schema.ContentBlockTypeServerToolCall {
// 模型已经消费了搜索结果——你只会看到后续的 reasoning/text block
fmt.Println("服务端搜索完成:", block.ServerToolCall.Name)
}
}

🧠 总结:上面六个坑覆盖了 90% 的 Eino 开发问题。记住三条原则:

  1. 打开要关(StreamReader → defer Close,AsyncIterator → gen.Close)
  2. 中断要透传(InterruptRerunError 不能被 Middleware 吞掉)
  3. 恢复要对齐(CheckPointStore + CheckPointID 跨 Runner 必须一致)

📂 源码仓库github.com/cloudwego/eino-examples

📖 继续阅读


感谢您的阅读!如果可以,给俺点些关注吧~

Eino学习笔记

周一 7月 20 2026
7722 · 46 分钟
封面
示例歌曲
示例艺术家
封面
示例歌曲
示例艺术家
0:00 / 0:00