vault backup: 2026-04-29 19:01:34
This commit is contained in:
@@ -1,24 +1,210 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-16"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第七章:Interrupt/Resume(中断与恢复)
|
||||
weight: 7
|
||||
tags: ["Eino", "Agent", "Interrupt", "Resume", "Backend", "DeepAgent", "审批流"]
|
||||
create time: "2026-04-29 15:30"
|
||||
---
|
||||
|
||||
本章目标:理解 Interrupt/Resume 机制,实现 Tool 审批流程,让用户在敏感操作前进行确认。
|
||||
# 第七章:Interrupt / Resume(中断与恢复)
|
||||
|
||||
## 代码位置
|
||||
## 概述
|
||||
|
||||
- 入口代码:[cmd/ch07/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch07/main.go)
|
||||
本章引入 Eino 的 **Interrupt / Resume** 机制——一种在人机协作中实现人工审批的能力。当 Agent 需要执行敏感操作(如删除文件、发送邮件、执行命令)时,可以在执行前暂停并等待用户确认;确认后继续,拒绝则返回错误。这是让 Agent 从"全自动"走向"安全可控"的关键一步。
|
||||
|
||||
## 前置条件
|
||||
## 为什么需要 Interrupt
|
||||
|
||||
与第一章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)。同时,需要与第四章一样设置 `PROJECT_ROOT`:
|
||||
前三章我们逐步为 Agent 添加工具能力,使其能够读取文件、搜索代码、执行命令。但全自动执行工具也存在风险:
|
||||
|
||||
```bash
|
||||
export PROJECT_ROOT=/path/to/eino # Eino 核心库根目录(不设置则默认使用当前目录)
|
||||
| 风险场景 | 后果 |
|
||||
|---------|------|
|
||||
| 误删文件 | 不可逆的数据丢失 |
|
||||
| 发送错误邮件 | 严重的沟通事故 |
|
||||
| 执行危险命令 | 系统环境被破坏 |
|
||||
| 修改关键配置 | 服务不可用 |
|
||||
|
||||
**Interrupt 的定位:**
|
||||
- **Interrupt 是 Agent 的暂停机制**:在关键操作前暂停,等待用户确认
|
||||
- **Interrupt 可携带信息**:向用户展示即将执行的操作详情
|
||||
- **Interrupt 可恢复**:确认后继续执行,拒绝后优雅返回错误
|
||||
|
||||
> [!tip] 简单类比
|
||||
>
|
||||
> - **自动执行** = "自动驾驶"(完全信任系统)
|
||||
> - **Interrupt** = "人工接管"(关键决策由人来做)
|
||||
|
||||
## 关键概念
|
||||
|
||||
### Interrupt 的两阶段执行
|
||||
|
||||
一个受审批保护的 Tool 在执行时被分成两个阶段:
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
A["Agent\n决定调用 Tool"] --> B{"Tool 内部"}
|
||||
B -->|"第一阶段"| C["保存参数"]
|
||||
C --> D["触发 Interrupt"]
|
||||
D --> E["Runner 暂停"]
|
||||
E --> F["向调用方返回\nInterrupt 事件"]
|
||||
F --> G["用户看到审批提示"]
|
||||
G --> H{"用户选择"}
|
||||
H -->|"批准"| I["runner.ResumeWith... 带上审批结果"]
|
||||
H -->|"拒绝"| J["runner.ResumeWith... 带上拒绝结果"]
|
||||
I --> K{"Tool 内部"}
|
||||
J --> K
|
||||
K -->|"第二阶段 Resume"| L["读取审批结果"]
|
||||
L -->|"Approved"| M["执行实际操作"]
|
||||
L -->|"Rejected"| N["操作被拒绝"]
|
||||
```
|
||||
|
||||
核心 API:
|
||||
|
||||
| API | 作用 |
|
||||
|------|------|
|
||||
| `tool.GetInterruptState[T](ctx)` | 判断当前是第一阶段还是 Resume 后的第二阶段 |
|
||||
| `tool.StatefulInterrupt(ctx, info, state)` | 触发中断,`info` 展示给用户,`state` 供 Resume 后取回 |
|
||||
| `tool.GetResumeContext[T](ctx)` | 获取用户的审批结果数据 |
|
||||
|
||||
> [!note] 两阶段设计精妙之处
|
||||
>
|
||||
> 同一个 Tool 函数被调用两次,通过 `GetInterruptState` 区分:第一次返回 false(触发中断),第二次返回 true(Resume 恢复)。这种"自反式"设计无需引入额外的状态机或外部协调器,中断逻辑就内聚在 Tool 自身内部。
|
||||
|
||||
### ApprovalMiddleware
|
||||
|
||||
生产实践中,推荐将中断逻辑放入 **Middleware** 而非每个 Tool 内部实现。这样审批规则集中管理、Tool 本身保持干净:
|
||||
|
||||
ApprovalMiddleware 拦截特定的 Tool 调用(如 `execute`),对每次调用统一施加审批逻辑:
|
||||
|
||||
```go
|
||||
type approvalMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *approvalMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
tCtx *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
// 仅拦截需审批的 Tool(例如 execute)
|
||||
if tCtx.Name != "execute" {
|
||||
return endpoint, nil
|
||||
}
|
||||
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
|
||||
|
||||
if !wasInterrupted {
|
||||
// 第一次调用 → 触发中断
|
||||
return "", tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: args,
|
||||
}, args)
|
||||
}
|
||||
|
||||
// Resume 阶段 → 检查用户是否批准
|
||||
isTarget, hasData, data := tool.GetResumeContext[*commontool.ApprovalResult](ctx)
|
||||
if isTarget && hasData {
|
||||
if data.Approved {
|
||||
return endpoint(ctx, storedArgs, opts...) // 通过中间件继续原 Tool 的执行
|
||||
}
|
||||
reason := ""
|
||||
if data.DisapproveReason != nil {
|
||||
reason = fmt.Sprintf(": %s", *data.DisapproveReason)
|
||||
}
|
||||
return fmt.Sprintf("tool '%s' disapproved%s", tCtx.Name, reason), nil
|
||||
}
|
||||
|
||||
// 非目标 Tool → 重新中断
|
||||
return "", tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: storedArgs,
|
||||
}, storedArgs)
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
> [!warning] Streamable 变体不可遗漏
|
||||
>
|
||||
> 如果 Agent 启用了流式输出(EnableStreaming: true),某些 Tool 调用可能走 `StreamableToolCall` 路径。此时必须同时实现 `WrapStreamableToolCall`,否则审批逻辑会被绕过。ch07 完整代码中两者都已覆盖。
|
||||
|
||||
### CheckPointStore
|
||||
|
||||
中断恢复还需要一个持久化组件来保存执行状态——这就是 `CheckPointStore`:
|
||||
|
||||
```go
|
||||
type CheckPointStore interface {
|
||||
Put(ctx context.Context, key string, checkpoint *Checkpoint) error
|
||||
Get(ctx context.Context, key string) (*Checkpoint, error)
|
||||
}
|
||||
```
|
||||
|
||||
它的作用不止于存储 Tool 参数,还包括 Runner 当前的执行进度。有了它,即使进程重启也能从中断点继续:
|
||||
|
||||
> [!example] CheckPointStore 的两种典型实现
|
||||
>
|
||||
> | 实现方式 | 适用场景 | 跨进程恢复 |
|
||||
> |---------|---------|----------|
|
||||
> | `adkstore.NewInMemoryStore()` | 开发调试、单进程 | ❌ |
|
||||
> | Redis / SQLite 等外部存储 | 生产部署 | ✅ |
|
||||
|
||||
## 代码实现
|
||||
|
||||
### 1. 配置 Runner 使用 CheckPointStore
|
||||
|
||||
```go
|
||||
runner := adk.NewRunner(ctx, adk.RunnerConfig{
|
||||
Agent: agent,
|
||||
EnableStreaming: true,
|
||||
CheckPointStore: adkstore.NewInMemoryStore(), // 内存存储
|
||||
})
|
||||
```
|
||||
|
||||
### 2. 配置 Agent 注册中间件
|
||||
|
||||
```go
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
// ... 其他配置
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&approvalMiddleware{}, // 审批中间件
|
||||
&safeToolMiddleware{}, // 将 Tool 错误转为字符串(中断类错误继续向上抛出)
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
### 3. 处理 Runner 返回的事件
|
||||
|
||||
```go
|
||||
checkPointID := sessionID
|
||||
|
||||
events := runner.Run(ctx, history, adk.WithCheckPointID(checkPointID))
|
||||
content, interruptInfo, err := printAndCollectAssistantFromEvents(events)
|
||||
|
||||
if interruptInfo != nil {
|
||||
// 使用同一个 stdin reader 读取「用户输入」与「审批 y/n」
|
||||
// 避免审批输入被误认为下一轮对话消息
|
||||
content, err = handleInterrupt(ctx, runner, checkPointID, interruptInfo, reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### 4. 完整的审批交互流程
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
U["用户:执行命令 echo hello"] --> S1["你> 请执行命令 echo hello"]
|
||||
S1 --> AGT["Runner.Run() 启动执行"]
|
||||
AGT --> A["Agent 分析意图\n决定调用 execute 工具"]
|
||||
A --> AM["ApprovalMiddleware\n拦截 Tool 调用"]
|
||||
AM --> SI["触发 StatefulInterrupt\n保存参数到 Store"]
|
||||
SI --> EVT["返回 Interrupt 事件"]
|
||||
EVT --> UI["控制台显示审批提示"]
|
||||
UI --> USER{"用户选择"}
|
||||
USER -->|"y"| RESUME["runner.ResumeWith\n携带审批结果 Approved=true"]
|
||||
USER -->|"n"| REJECT_DIRECT["runner.ResumeWith\n携带审批结果 Approved=false"]
|
||||
RESUME --> RTOOL["Tool 再次被调用\nGetInterruptState = true\n读取审批结果并批准"]
|
||||
RTOOL --> EXEC["执行 execute\necho hello"]
|
||||
EXEC --> OUT["输出: hello"]
|
||||
REJECT_DIRECT --> RTOOL2["Tool 再次被调用\nGetInterruptState = true\n读取审批结果并拒绝"]
|
||||
RTOOL2 --> NOP["输出: tool disapproved"]
|
||||
```
|
||||
|
||||
## 运行
|
||||
@@ -26,9 +212,7 @@ export PROJECT_ROOT=/path/to/eino # Eino 核心库根目录(不设置则默
|
||||
在 `examples/quickstart/chatwitheino` 目录下执行:
|
||||
|
||||
```bash
|
||||
# 设置项目根目录
|
||||
export PROJECT_ROOT=/path/to/your/project
|
||||
|
||||
export PROJECT_ROOT=/path/to/eino # Eino 核心库根目录(不设置则默认使用当前目录)
|
||||
go run ./cmd/ch07
|
||||
```
|
||||
|
||||
@@ -47,309 +231,42 @@ Approve this action? (y/n): y
|
||||
hello
|
||||
```
|
||||
|
||||
## 从自动执行到人工审批:为什么需要 Interrupt
|
||||
|
||||
前几章我们实现的 Agent 会自动执行所有 Tool 调用,但在某些场景下这是危险的:
|
||||
|
||||
**自动执行的风险:**
|
||||
|
||||
- 删除文件:误删重要数据
|
||||
- 发送邮件:发送错误内容
|
||||
- 执行命令:执行危险操作
|
||||
- 修改配置:破坏系统设置
|
||||
|
||||
**Interrupt 的定位:**
|
||||
|
||||
- **Interrupt 是 Agent 的暂停机制**:在关键操作前暂停,等待用户确认
|
||||
- **Interrupt 可携带信息**:向用户展示即将执行的操作
|
||||
- **Interrupt 可恢复**:用户确认后继续执行,拒绝后返回错误
|
||||
|
||||
**简单类比:**
|
||||
|
||||
- **自动执行** = "自动驾驶"(完全信任系统)
|
||||
- **Interrupt** = "人工接管"(关键决策由人来做)
|
||||
|
||||
## 关键概念
|
||||
|
||||
### Interrupt 机制
|
||||
|
||||
`Interrupt` 是 Eino 中实现人机协作的核心机制。
|
||||
|
||||
**核心思想:在执行关键操作前暂停,等待用户确认后继续。**
|
||||
|
||||
一个需要审批的 Tool 的执行被分成**两个阶段**:
|
||||
|
||||
1. **第一次调用(触发中断)**:Tool 保存当前参数,然后返回一个中断信号。Runner 暂停执行,向调用侧返回 Interrupt 事件。
|
||||
2. **用户审批后恢复(Resume)**:Runner 重新调用 Tool,此时 Tool 检测到"已中断过",直接读取用户的审批结果并执行(或拒绝)。
|
||||
|
||||
**简化版伪代码:**
|
||||
|
||||
```
|
||||
func myTool(ctx, args):
|
||||
if 第一次调用:
|
||||
保存 args
|
||||
return 中断信号 // Runner 暂停,展示审批提示
|
||||
else: // Resume 后的第二次调用
|
||||
if 用户批准:
|
||||
return 执行操作(保存的 args)
|
||||
else:
|
||||
return "操作被用户拒绝"
|
||||
```
|
||||
|
||||
**完整代码及关键字段说明:**
|
||||
|
||||
```go
|
||||
// 在 Tool 中触发中断
|
||||
func myTool(ctx context.Context, args string) (string, error) {
|
||||
// wasInterrupted: 是否是 Resume 后的第二次调用(第一次为 false,Resume 后为 true)
|
||||
// storedArgs: 第一次调用时通过 StatefulInterrupt 保存的参数,Resume 后可取回
|
||||
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
|
||||
|
||||
if !wasInterrupted {
|
||||
// 第一次调用:触发中断,同时保存 args 供 Resume 后使用
|
||||
return "", tool.StatefulInterrupt(ctx, &ApprovalInfo{
|
||||
ToolName: "my_tool",
|
||||
ArgumentsInJSON: args,
|
||||
}, args) // 第三个参数是要保存的状态(Resume 后通过 storedArgs 取回)
|
||||
}
|
||||
|
||||
// Resume 后的第二次调用:读取用户审批结果
|
||||
// isTarget: 本次 Resume 是否针对当前 Tool(一次 Resume 只针对一个 Tool)
|
||||
// hasData: Resume 时是否携带了审批结果数据
|
||||
// data: 用户传入的审批结果
|
||||
isTarget, hasData, data := tool.GetResumeContext[*ApprovalResult](ctx)
|
||||
if isTarget && hasData {
|
||||
if data.Approved {
|
||||
return doSomething(storedArgs) // 使用保存的参数执行实际操作
|
||||
}
|
||||
return "Operation rejected by user", nil
|
||||
}
|
||||
|
||||
// 其他情况(isTarget=false 意味着本次 Resume 目标不是当前 Tool):重新中断
|
||||
return "", tool.StatefulInterrupt(ctx, &ApprovalInfo{
|
||||
ToolName: "my_tool",
|
||||
ArgumentsInJSON: storedArgs,
|
||||
}, storedArgs)
|
||||
}
|
||||
```
|
||||
|
||||
### ApprovalMiddleware
|
||||
|
||||
`ApprovalMiddleware` 是一个通用的审批中间件,可以拦截特定 Tool 的调用:
|
||||
|
||||
```go
|
||||
type approvalMiddleware struct {
|
||||
*adk.BaseChatModelAgentMiddleware
|
||||
}
|
||||
|
||||
func (m *approvalMiddleware) WrapInvokableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.InvokableToolCallEndpoint,
|
||||
tCtx *adk.ToolContext,
|
||||
) (adk.InvokableToolCallEndpoint, error) {
|
||||
// 只拦截需要审批的 Tool
|
||||
if tCtx.Name != "execute" {
|
||||
return endpoint, nil
|
||||
}
|
||||
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
|
||||
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
|
||||
|
||||
if !wasInterrupted {
|
||||
return "", tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: args,
|
||||
}, args)
|
||||
}
|
||||
|
||||
isTarget, hasData, data := tool.GetResumeContext[*commontool.ApprovalResult](ctx)
|
||||
if isTarget && hasData {
|
||||
if data.Approved {
|
||||
return endpoint(ctx, storedArgs, opts...)
|
||||
}
|
||||
if data.DisapproveReason != nil {
|
||||
return fmt.Sprintf("tool '%s' disapproved: %s", tCtx.Name, *data.DisapproveReason), nil
|
||||
}
|
||||
return fmt.Sprintf("tool '%s' disapproved", tCtx.Name), nil
|
||||
}
|
||||
|
||||
// 重新中断
|
||||
return "", tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: storedArgs,
|
||||
}, storedArgs)
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (m *approvalMiddleware) WrapStreamableToolCall(
|
||||
_ context.Context,
|
||||
endpoint adk.StreamableToolCallEndpoint,
|
||||
tCtx *adk.ToolContext,
|
||||
) (adk.StreamableToolCallEndpoint, error) {
|
||||
// 如果 agent 配置了 StreamingShell,则 execute 会走流式调用,需要实现该方法才能拦截到
|
||||
if tCtx.Name != "execute" {
|
||||
return endpoint, nil
|
||||
}
|
||||
return func(ctx context.Context, args string, opts ...tool.Option) (*schema.StreamReader[string], error) {
|
||||
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
|
||||
if !wasInterrupted {
|
||||
return nil, tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: args,
|
||||
}, args)
|
||||
}
|
||||
|
||||
isTarget, hasData, data := tool.GetResumeContext[*commontool.ApprovalResult](ctx)
|
||||
if isTarget && hasData {
|
||||
if data.Approved {
|
||||
return endpoint(ctx, storedArgs, opts...)
|
||||
}
|
||||
if data.DisapproveReason != nil {
|
||||
return singleChunkReader(fmt.Sprintf("tool '%s' disapproved: %s", tCtx.Name, *data.DisapproveReason)), nil
|
||||
}
|
||||
return singleChunkReader(fmt.Sprintf("tool '%s' disapproved", tCtx.Name)), nil
|
||||
}
|
||||
|
||||
isTarget, _, _ = tool.GetResumeContext[any](ctx)
|
||||
if !isTarget {
|
||||
return nil, tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: tCtx.Name,
|
||||
ArgumentsInJSON: storedArgs,
|
||||
}, storedArgs)
|
||||
}
|
||||
|
||||
return endpoint(ctx, storedArgs, opts...)
|
||||
}, nil
|
||||
}
|
||||
```
|
||||
|
||||
### CheckPointStore
|
||||
|
||||
`CheckPointStore` 是实现中断恢复的关键组件:
|
||||
|
||||
```go
|
||||
type CheckPointStore interface {
|
||||
// 保存检查点
|
||||
Put(ctx context.Context, key string, checkpoint *Checkpoint) error
|
||||
|
||||
// 获取检查点
|
||||
Get(ctx context.Context, key string) (*Checkpoint, error)
|
||||
}
|
||||
```
|
||||
|
||||
**为什么需要 CheckPointStore?**
|
||||
|
||||
- 中断时保存状态:Tool 参数、执行位置等
|
||||
- 恢复时加载状态:从中断点继续执行
|
||||
- 支持跨进程恢复:进程重启后仍可恢复
|
||||
|
||||
## Interrupt/Resume 的实现
|
||||
|
||||
### 1. 配置 Runner 使用 CheckPointStore
|
||||
|
||||
```go
|
||||
runner := adk.NewRunner(ctx, adk.RunnerConfig{
|
||||
Agent: agent,
|
||||
EnableStreaming: true,
|
||||
CheckPointStore: adkstore.NewInMemoryStore(), // 内存存储
|
||||
})
|
||||
```
|
||||
|
||||
### 2. 配置 Agent 使用 ApprovalMiddleware
|
||||
|
||||
```go
|
||||
agent, err := deep.New(ctx, &deep.Config{
|
||||
// ... 其他配置
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
&approvalMiddleware{}, // 添加审批中间件
|
||||
&safeToolMiddleware{}, // 将 Tool 错误转换为字符串(中断类错误会继续向上抛出)
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
### 3. 处理中断事件
|
||||
|
||||
```go
|
||||
checkPointID := sessionID
|
||||
|
||||
events := runner.Run(ctx, history, adk.WithCheckPointID(checkPointID))
|
||||
content, interruptInfo, err := printAndCollectAssistantFromEvents(events)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if interruptInfo != nil {
|
||||
// 注意:建议使用同一个 stdin reader 同时读取「用户输入」与「审批 y/n」
|
||||
// 避免审批输入被当成下一轮 you> 的消息
|
||||
content, err = handleInterrupt(ctx, runner, checkPointID, interruptInfo, reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
_ = session.Append(schema.AssistantMessage(content, nil))
|
||||
```
|
||||
|
||||
## Interrupt/Resume 执行流程
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 用户:执行命令 echo hello │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 分析意图 │
|
||||
│ 决定调用 execute │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ ApprovalMiddleware │
|
||||
│ 拦截 Tool 调用 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 触发 Interrupt │
|
||||
│ 保存状态到 Store │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 返回 Interrupt 事件 │
|
||||
│ 等待用户审批 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 用户输入 y/n │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ runner.ResumeWith... │
|
||||
│ 恢复执行 │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 执行 execute │
|
||||
│ 或返回拒绝信息 │
|
||||
└──────────────────────┘
|
||||
```
|
||||
> [!question] 深入思考
|
||||
>
|
||||
> 上面的输出中有两条 `hello`——一条来自 `[tool result]`,另一条是 Assistant 的最终回复。你能解释它们分别来自哪里吗?
|
||||
> 提示:第一条是 `printAndCollectAssistantFromEvents` 对流式事件中 Tool Result 片段的打印,第二条是 Agent 整合信息后生成的自然语言回复。理解了这一点,你就掌握了 Eino 事件模型的核心。
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Interrupt**:Agent 的暂停机制,在关键操作前暂停等待确认
|
||||
- **Resume**:恢复执行,用户确认后继续或拒绝后返回错误
|
||||
- **ApprovalMiddleware**:通用审批中间件,拦截特定 Tool 调用
|
||||
- **CheckPointStore**:保存中断状态,支持跨进程恢复
|
||||
- **人机协作**:关键决策由人来确认,提高安全性
|
||||
| 概念 | 一句话理解 |
|
||||
|------|-----------|
|
||||
| **Interrupt** | Agent 在敏感操作前的暂停机制 |
|
||||
| **Resume** | 用户审批后恢复执行,支持批准与拒绝两种结果 |
|
||||
| **Two-stage Execution** | 同一个 Tool 被调用两次,通过 `GetInterruptState` 区分阶段 |
|
||||
| **ApprovalMiddleware** | 集中式拦截特定 Tool 的审批逻辑,使 Tool 保持干净 |
|
||||
| **CheckPointStore** | 保存中断状态和执行位置,支持跨进程恢复 |
|
||||
| **人机协作** | 关键决策由人类确认,兼顾 Agent 自动化与安全可控 |
|
||||
|
||||
## 扩展思考
|
||||
|
||||
**其他 Interrupt 场景:**
|
||||
### 更多 Interrupt 应用场景
|
||||
|
||||
- 多选项审批:用户选择多个选项之一
|
||||
- 参数补全:用户提供缺失的参数
|
||||
- 条件分支:用户决定执行路径
|
||||
| 场景 | 说明 |
|
||||
|------|------|
|
||||
| 多选项审批 | 用户从多个选项中选择一个(而非简单的 y/n) |
|
||||
| 参数补全 | 用户提供缺失的参数值后才继续执行 |
|
||||
| 条件分支 | 用户决定不同的执行路径 |
|
||||
|
||||
**审批策略:**
|
||||
### 审批策略
|
||||
|
||||
- 白名单:只审批敏感操作
|
||||
- 黑名单:审批所有操作,除了安全的
|
||||
- 动态规则:根据参数内容决定是否审批
|
||||
| 策略 | 适用场景 |
|
||||
|------|---------|
|
||||
| 白名单 | 只审批极少数敏感操作(推荐默认做法) |
|
||||
| 黑名单 | 审批所有操作,除已知的安全操作外 |
|
||||
| 动态规则 | 根据参数内容决定是否审批(如文件大小、操作范围) |
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/_index]]
|
||||
- [[Eino/quick_start/chapter_04_tool_and_filesystem]] — 文件系统访问与 DeepAgent(第三章的工具章节)
|
||||
- [[Eino/quick_start/chapter_05_middleware]] — Middleware 模式详解(上一章)
|
||||
|
||||
@@ -1,13 +1,17 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-12"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第八章:Graph Tool(复杂工作流)
|
||||
weight: 8
|
||||
tags: ["Eino", "Agent", "GraphTool", "Compose", "Workflow", "Backend"]
|
||||
create time: "2026-04-29 15:30"
|
||||
---
|
||||
|
||||
本章目标:理解 Graph Tool 的概念,实现大文件的并行 chunk 召回,引入 compose 包构建复杂工作流。
|
||||
# 第八章:Graph Tool(复杂工作流)
|
||||
|
||||
## 概述
|
||||
|
||||
本章引入 Eino 的 **Graph Tool** 能力——将复杂的编排工作流封装为一个可调用的 Tool。通过 `compose.Workflow` 构建包含读取、分块、并行评分、筛选和答案生成的多步骤流水线,让 Agent 能够处理需要多阶段协同的大文件 RAG 场景。
|
||||
|
||||
> [!tip] 一句话理解 Graph Tool
|
||||
>
|
||||
> **简单 Tool = 单步操作**(如读取文件),**Graph Tool = 完整流水线**(读取 → 分块 → 并行评分 → 筛选 → 生成答案)。它是 compose 编排能力的 Tool 化封装入口。
|
||||
|
||||
## 代码位置
|
||||
|
||||
@@ -45,94 +49,119 @@ you> 请帮我分析 RFC6455 文档中关于 WebSocket 握手的部分
|
||||
|
||||
**简单 Tool 的局限:**
|
||||
|
||||
- 单一职责:每个 Tool 只做一件事
|
||||
- 无法并行:多个独立任务无法同时执行
|
||||
- 难以复用:复杂逻辑难以拆分和组合
|
||||
| 局限 | 说明 |
|
||||
|------|------|
|
||||
| 单一职责 | 每个 Tool 只能做一件事(读文件、搜索等) |
|
||||
| 无法并行 | 多个独立子任务不能同时执行 |
|
||||
| 难以复用 | 复杂逻辑硬编码在调用链中,无法单独测试和复用 |
|
||||
|
||||
**重要说明:本章只是展示 compose/graph/workflow 能力的一角。**
|
||||
|
||||
从更大的视角看,Eino 的 `compose` 包提供了非常通用、确定性的编排能力:你可以把任何需要"确定性业务流程"的系统,用 `compose` 的 Graph/Chain/Workflow 组织成可执行的流水线,并且它能够**原生编排 Eino 的所有 component**(如 ChatModel、Prompt、Tools、Retriever、Embedding、Indexer 等),同时具备完整的 **callback** 体系,以及 **interrupt/resume + checkpoint** 支持。
|
||||
从更大的视角看,Eino 的 `compose` 包提供了非常通用、确定性的编排能力:你可以把任何需要"确定性业务流程"的系统,用 `compose` 的 Graph/Chain/Workflow 组织成可执行的流水线,并且它能够**原生编排 Eino 的所有 component**(ChatModel、Prompt、Tools、Retriever、Embedding、Indexer 等),同时具备完整的 **callback** 体系,以及 **interrupt/resume + checkpoint** 支持。
|
||||
|
||||
**Graph Tool 的定位:**
|
||||
### Graph Tool 的定位
|
||||
|
||||
- **Graph Tool 是 compose 工作流的 Tool 化封装**:把 `compose.Graph / compose.Chain / compose.Workflow` 这类可编译的编排产物,包装成一个 Agent 可调用的 Tool
|
||||
- **支持并行/分支/组合**:由 compose 提供(并行、分支、字段映射、子图等),Graph Tool 只是把它们暴露为 Tool 入口
|
||||
- **支持状态管理与持久化**:节点间传递数据、以及通过 checkpoint 保存/恢复运行状态
|
||||
- **可中断恢复**:既支持工作流内部的中断(节点里触发 interrupt),也支持工具层面的中断包装(嵌套 interrupt 场景)
|
||||
> [!note] Graph Tool vs 简单 Tool
|
||||
>
|
||||
> | 对比项 | 简单 Tool | Graph Tool |
|
||||
> |--------|----------|------------|
|
||||
> | 本质 | 单步函数 | compose 编排产物的封装 |
|
||||
> | 编排 | 无 | 由 compose 提供(并行、分支、字段映射) |
|
||||
> | 状态管理 | 无 | 节点间传递数据 + checkpoint 持久化 |
|
||||
> | 中断恢复 | 不支持 | 支持(嵌套 interrupt 场景) |
|
||||
|
||||
**简单类比:**
|
||||
### 核心类比
|
||||
|
||||
- **简单 Tool** = "单步操作"(读取文件)
|
||||
- **Graph Tool** = "流水线"(读取 → 分块 → 评分 → 筛选 → 生成答案)
|
||||
> [!tip] 厨房做菜类比
|
||||
>
|
||||
> - **简单 Tool**:像是一个厨具(菜刀——只负责切东西)
|
||||
> - **Graph Tool**:像是一条预制菜流水线(备料 → 烹饪 → 摆盘——每一步自动衔接,你只需说"做这道菜")
|
||||
|
||||
## 关键概念
|
||||
|
||||
### compose.Workflow
|
||||
|
||||
`compose.Workflow` 是 Eino 中构建工作流的核心组件:
|
||||
`compose.Workflow` 是 Eino 中构建有状态工作流的核心组件。与线性 Chain 不同,Workflow 允许创建 DAG(有向无环图),支持汇聚节点、并行分支和非相邻连接:
|
||||
|
||||
```go
|
||||
wf := compose.NewWorkflow[Input, Output]()
|
||||
|
||||
// 添加节点
|
||||
// 添加节点并建立连接
|
||||
wf.AddLambdaNode("load", loadFunc).AddInput(compose.START)
|
||||
wf.AddLambdaNode("chunk", chunkFunc).AddInput("load")
|
||||
wf.AddLambdaNode("score", scoreFunc).AddInput("chunk")
|
||||
wf.AddLambdaNode("answer", answerFunc).AddInput("score")
|
||||
|
||||
// 连接到结束节点
|
||||
wf.AddLambdaNode("answer", answerFunc).
|
||||
AddInput("chunk").
|
||||
AddInputWithOptions(compose.START,
|
||||
[]*compose.FieldMapping{compose.MapFields("Question", "Question")},
|
||||
compose.WithNoDirectDependency())
|
||||
wf.End().AddInput("answer")
|
||||
```
|
||||
|
||||
**核心概念:**
|
||||
> [!question] 深入思考
|
||||
>
|
||||
> Workflow 为什么需要 `START` 和 `END` 这两个虚拟节点,而不是直接指定输入输出?
|
||||
> 提示:想想如果工作流有多个入口点(例如用户可以直接跳转到某个中间节点重试),或者需要在运行时动态插入新节点。START/END 为这些灵活性提供了统一的锚点。
|
||||
|
||||
- **Node**:工作流中的处理单元
|
||||
- **Edge**:节点间的数据流向
|
||||
- **START**:工作流的入口
|
||||
- **END**:工作流的出口
|
||||
### BatchNode(并行处理)
|
||||
|
||||
### BatchNode
|
||||
|
||||
`BatchNode` 用于并行处理多个任务:
|
||||
`BatchNode` 用于并行处理一批独立任务,充分利用计算资源:
|
||||
|
||||
```go
|
||||
scorer := batch.NewBatchNode(&batch.NodeConfig[Task, Result]{
|
||||
scorer := batch.NewBatchNode(&batch.NodeConfig[scoreTask, scoredChunk]{
|
||||
Name: "ChunkScorer",
|
||||
InnerTask: scoreOneChunk, // 单个任务的处理函数
|
||||
MaxConcurrency: 5, // 最大并发数
|
||||
InnerTask: newScoreWorkflow(cm), // 单个 chunk 的评分流程
|
||||
MaxConcurrency: 5, // 最大并发数
|
||||
})
|
||||
```
|
||||
|
||||
**工作原理:**
|
||||
|
||||
1. 接收任务列表作为输入
|
||||
2. 并行执行每个任务(受 MaxConcurrency 限制)
|
||||
3. 收集所有结果返回
|
||||
1. 接收任务切片作为输入
|
||||
2. 按 `MaxConcurrency` 限制并行调度(内部使用 goroutine pool)
|
||||
3. 所有结果收集后按顺序返回
|
||||
|
||||
### FieldMapping
|
||||
> [!tip] 选择 MaxConcurrency 的原则
|
||||
>
|
||||
> - 过低 → 浪费了并发能力,响应慢
|
||||
> - 过高 → 资源竞争,LLM API 限流
|
||||
> - 推荐做法:以 LLM Provider 的 QPS 上限为参考值,一般 3~10 之间调整
|
||||
|
||||
`FieldMapping` 用于跨节点传递数据:
|
||||
### FieldMapping(跨节点数据传递)
|
||||
|
||||
FieldMapping 解决非相邻节点间的数据传递问题:当两个节点没有直接的边连接时,你需要显式声明数据的来源和目标字段。
|
||||
|
||||
```go
|
||||
wf.AddLambdaNode("answer", answerFunc).
|
||||
AddInputWithOptions("filter", // 从 filter 节点获取数据
|
||||
[]*compose.FieldMapping{compose.ToField("TopK")},
|
||||
wf.AddLambdaNode("score", scoreFunc).
|
||||
// 从 "chunk" 节点取 All 数据,映射到当前节点的 Chunks 字段
|
||||
AddInputWithOptions("chunk",
|
||||
[]*compose.FieldMapping{compose.ToField("Chunks")},
|
||||
compose.WithNoDirectDependency()).
|
||||
AddInputWithOptions(compose.START, // 从 START 节点获取数据
|
||||
// 从 START 节点取 Question 字段,直接映射到当前节点的 Question 字段
|
||||
AddInputWithOptions(compose.START,
|
||||
[]*compose.FieldMapping{compose.MapFields("Question", "Question")},
|
||||
compose.WithNoDirectDependency())
|
||||
```
|
||||
|
||||
**为什么需要 FieldMapping?**
|
||||
**三种 FieldMapping 方式:**
|
||||
|
||||
- 非相邻节点间传递数据
|
||||
- 多个数据源合并到同一节点
|
||||
- 数据字段重命名
|
||||
| 方法 | 作用 | 适用场景 |
|
||||
|------|------|---------|
|
||||
| `MapFields(src, dst)` | 字段重命名映射 | 两端字段名不一致时 |
|
||||
| `ToField(dst)` | 整条数据映射到单一字段 | 上游只有一个输出,且需包裹到 struct |
|
||||
| `All()` | 传入上游全部输出(默认行为) | 相邻节点间的直接传递 |
|
||||
|
||||
**为什么非相邻节点需要 `WithNoDirectDependency`?**
|
||||
|
||||
Eino 依赖图检测会验证节点的输入是否来自前驱节点。当使用 FieldMapping 跨越层级取值时,必须显式标记 `WithNoDirectDependency()`,否则会被依赖检查拦截。
|
||||
|
||||
## Graph Tool 的实现
|
||||
|
||||
下面我们以"大文件内容检索并回答"为例,逐步构建一个完整的 Graph Tool。整个流程分为三步:定义 IO 结构 → 构建工作流 → 封装为 Tool。
|
||||
|
||||
### 1. 定义输入输出结构
|
||||
|
||||
输入和输出定义了 Graph Tool 对外暴露的接口契约,也是 Agent 调用时的参数 schema 来源:
|
||||
|
||||
```go
|
||||
type Input struct {
|
||||
FilePath string `json:"file_path" jsonschema:"description=Absolute path to the document"`
|
||||
@@ -145,13 +174,19 @@ type Output struct {
|
||||
}
|
||||
```
|
||||
|
||||
> [!note] jsonschema tag 的作用
|
||||
>
|
||||
> 这些标签会被自动转换为 JSON Schema,决定了 Agent(LLM)看到的工具参数描述。写得好,模型就能精准理解该传什么值。
|
||||
|
||||
### 2. 构建工作流
|
||||
|
||||
完整的 `buildWorkflow` 函数实现了五个阶段的流水线:
|
||||
|
||||
```go
|
||||
func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
wf := compose.NewWorkflow[Input, Output]()
|
||||
|
||||
// load: 读取文件
|
||||
// --- load: 读取文件 ---
|
||||
wf.AddLambdaNode("load", compose.InvokableLambda(
|
||||
func(ctx context.Context, in Input) ([]*schema.Document, error) {
|
||||
data, err := os.ReadFile(in.FilePath)
|
||||
@@ -162,7 +197,7 @@ func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
},
|
||||
)).AddInput(compose.START)
|
||||
|
||||
// chunk: 分块
|
||||
// --- chunk: 分块 ---
|
||||
wf.AddLambdaNode("chunk", compose.InvokableLambda(
|
||||
func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
|
||||
var out []*schema.Document
|
||||
@@ -173,7 +208,7 @@ func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
},
|
||||
)).AddInput("load")
|
||||
|
||||
// score: 并行评分
|
||||
// --- score: 并行评分(核心亮点)---
|
||||
scorer := batch.NewBatchNode(&batch.NodeConfig[scoreTask, scoredChunk]{
|
||||
Name: "ChunkScorer",
|
||||
InnerTask: newScoreWorkflow(cm),
|
||||
@@ -192,13 +227,12 @@ func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
AddInputWithOptions("chunk", []*compose.FieldMapping{compose.ToField("Chunks")}, compose.WithNoDirectDependency()).
|
||||
AddInputWithOptions(compose.START, []*compose.FieldMapping{compose.MapFields("Question", "Question")}, compose.WithNoDirectDependency())
|
||||
|
||||
// filter: 筛选 top-k
|
||||
// --- filter: 筛选 top-k ---
|
||||
wf.AddLambdaNode("filter", compose.InvokableLambda(
|
||||
func(ctx context.Context, scored []scoredChunk) ([]scoredChunk, error) {
|
||||
sort.Slice(scored, func(i, j int) bool {
|
||||
return scored[i].Score > scored[j].Score
|
||||
})
|
||||
// 返回 top-3
|
||||
if len(scored) > 3 {
|
||||
scored = scored[:3]
|
||||
}
|
||||
@@ -206,13 +240,8 @@ func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
},
|
||||
)).AddInput("score")
|
||||
|
||||
// answer: 生成答案
|
||||
wf.AddLambdaNode("answer", compose.InvokableLambda(
|
||||
func(ctx context.Context, in synthIn) (Output, error) {
|
||||
return synthesize(ctx, cm, in)
|
||||
},
|
||||
)).
|
||||
AddInputWithOptions("filter", []*compose.FieldMapping{compose.ToField("TopK")}, compose.WithNoDirectDependency()).
|
||||
// --- answer: 生成最终答案 ---
|
||||
wf.AddInputWithOptions("filter", []*compose.FieldMapping{compose.ToField("TopK")}, compose.WithNoDirectDependency()).
|
||||
AddInputWithOptions(compose.START, []*compose.FieldMapping{compose.MapFields("Question", "Question")}, compose.WithNoDirectDependency())
|
||||
|
||||
wf.End().AddInput("answer")
|
||||
@@ -221,93 +250,119 @@ func buildWorkflow(cm model.BaseChatModel) *compose.Workflow[Input, Output] {
|
||||
}
|
||||
```
|
||||
|
||||
> [!note] 代码解读:为什么 score 和 answer 都有两处 AddInput?
|
||||
>
|
||||
> **score 节点**需要两个数据来源:
|
||||
> - `chunk` 的输出(待评分的文本块)
|
||||
> - `START` 的 `Question`(用户的问题,用来给每个 block 打分)
|
||||
>
|
||||
> **answer 节点**同理也需要:
|
||||
> - `filter` 的输出(top-k 的相关片段)
|
||||
> - `START` 的 `Question`(拼接到 prompt 中)
|
||||
>
|
||||
> 这就是为什么需要 `WithNoDirectDependency()`——它们跳过了中间节点,直接向源头要数据。
|
||||
|
||||
### 3. 封装为 Tool
|
||||
|
||||
最后一步是将编译后的工作流包装成 Agent 可调用的标准 Tool:
|
||||
|
||||
```go
|
||||
func BuildTool(ctx context.Context, cm model.BaseChatModel) (tool.BaseTool, error) {
|
||||
wf := buildWorkflow(cm)
|
||||
return graphtool.NewInvokableGraphTool[Input, Output](
|
||||
wf,
|
||||
"answer_from_document",
|
||||
"Search a large document for relevant content and synthesize an answer.",
|
||||
"answer_from_document", // Tool 名称(Agent 看到的名字)
|
||||
"Search a large document for relevant content and synthesize an answer.", // Tool 描述
|
||||
)
|
||||
}
|
||||
```
|
||||
|
||||
**关键代码片段(**注意:这是简化后的代码片段,不能直接运行,完整代码请参考** [rag/rag.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/rag/rag.go)):
|
||||
> [!warning] 编译时机
|
||||
>
|
||||
> `graphtool.NewInvokableGraphTool` 内部会对 Workflow 执行编译检查,验证节点连通性、类型兼容性。如果在运行时才发现错误,排查会比较困难——建议在单元测试中对 buildWorkflow 的返回值做一次 compile-time check。
|
||||
|
||||
## Graph Tool 执行流程图
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["输入: file_path, question"] --> B["load\n读取文件\n→ []*Document"]
|
||||
B --> C["chunk\n分块 (800 tokens)\n→ []*Document"]
|
||||
|
||||
C --> D["score\n并行评分\n(MaxConcurrency=5)\n→ []scoredChunk"]
|
||||
|
||||
D --> E["filter\n排序并取 top-k\n→ []scoredChunk"]
|
||||
|
||||
E --> F["answer\n结合问题和\nTop-K 片段生成答案\n→ Output"]
|
||||
|
||||
F --> G["返回: {answer, sources}"]
|
||||
|
||||
style A fill:#e3f2fd
|
||||
style G fill:#e8f5e9
|
||||
style D fill:#fff3e0
|
||||
```
|
||||
|
||||
**流程中的关键设计决策:**
|
||||
|
||||
| 阶段 | 决策点 | 原因 |
|
||||
|------|--------|------|
|
||||
| chunk | 固定 800 token 分块 | 平衡上下文窗口与检索精度 |
|
||||
| score | 并行评分(MaxConcurrency=5) | 避免串行等待,利用 LLM API 并发能力 |
|
||||
| filter | 保留 top-3 | 控制后续 token 消耗,避免信息过载 |
|
||||
|
||||
## 可中断恢复
|
||||
|
||||
Graph Tool 天然继承 Eino 的中断恢复机制。当工作流内部的某个节点触发 `interrupt` 时,Runner 会暂停整个工作流,等待用户输入后 resume:
|
||||
|
||||
```go
|
||||
// 构建工作流
|
||||
wf := compose.NewWorkflow[Input, Output]()
|
||||
|
||||
// 添加节点
|
||||
wf.AddLambdaNode("load", loadFunc).AddInput(compose.START)
|
||||
wf.AddLambdaNode("chunk", chunkFunc).AddInput("load")
|
||||
wf.AddLambdaNode("score", scoreFunc).
|
||||
AddInputWithOptions("chunk", []*compose.FieldMapping{compose.ToField("Chunks")}, compose.WithNoDirectDependency()).
|
||||
AddInputWithOptions(compose.START, []*compose.FieldMapping{compose.MapFields("Question", "Question")}, compose.WithNoDirectDependency())
|
||||
|
||||
// 封装为 Tool
|
||||
return graphtool.NewInvokableGraphTool[Input, Output](wf, "answer_from_document", "...")
|
||||
// 在工作流节点中使用 interrupt
|
||||
func myNode(ctx context.Context, input MyInput) (MyOutput, error) {
|
||||
wasInterrupted, _, stored := tool.GetInterruptState[string](ctx)
|
||||
if !wasInterrupted {
|
||||
return MyOutput{}, tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
|
||||
ToolName: "my_workflow_step",
|
||||
ArgumentsInJSON: stored,
|
||||
}, stored)
|
||||
}
|
||||
// Resume 后继续执行...
|
||||
return process(stored), nil
|
||||
}
|
||||
```
|
||||
|
||||
## Graph Tool 执行流程
|
||||
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 输入:file_path, question │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ load: 读取文件 │
|
||||
│ 输出: []*Document │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ chunk: 分块 │
|
||||
│ 输出: []*Document │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ score: 并行评分 │
|
||||
│ (MaxConcurrency=5) │
|
||||
│ 输出: []scoredChunk │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ filter: 筛选 top-k │
|
||||
│ 输出: []scoredChunk │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ answer: 生成答案 │
|
||||
│ 输出: Output │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 返回结果 │
|
||||
│ {answer, sources} │
|
||||
└──────────────────────┘
|
||||
```
|
||||
> [!tip] Graph Tool 的中断优势
|
||||
>
|
||||
> 由于每个节点都是独立的 lambda 函数,可以在任意节点插入 interrupt 逻辑,而无需修改其他节点。这种细粒度的可控性是简单 Tool 无法做到的。
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **Graph Tool**:将复杂工作流封装为 Tool,支持多步骤协同
|
||||
- **compose.Workflow**:构建工作流的核心组件
|
||||
- **BatchNode**:并行处理多个任务
|
||||
- **FieldMapping**:跨节点传递数据
|
||||
- **可中断恢复**:Graph Tool 支持 Checkpoint 机制
|
||||
| 概念 | 一句话理解 |
|
||||
|------|-----------|
|
||||
| **Graph Tool** | 将 compose 编排产物封装为 Agent 可调用的 Tool 入口 |
|
||||
| **compose.Workflow** | 支持 DAG 结构的有状态工作流,可表达复杂业务逻辑 |
|
||||
| **BatchNode** | 并行处理批量任务的内置组件,受 MaxConcurrency 限制 |
|
||||
| **FieldMapping** | 跨节点传递数据的机制,解决非相邻节点间的通信 |
|
||||
| **可中断恢复** | Graph Tool 完整继承 interrupt/resume + checkpoint 能力 |
|
||||
|
||||
## 扩展思考
|
||||
|
||||
**其他 Graph Tool 应用:**
|
||||
### Graph Tool 的典型应用场景
|
||||
|
||||
- 多文档 RAG:并行处理多个文档
|
||||
- 多模型协作:不同模型处理不同任务
|
||||
- 复杂决策树:根据条件选择不同分支
|
||||
| 场景 | 说明 | 收益 |
|
||||
|------|------|------|
|
||||
| **多文档 RAG** | 并行检索多个文档源并综合回答 | 减少 Token 往返次数,一次 Tool Call 覆盖全部 |
|
||||
| **多模型协作** | 不同模型处理不同阶段(摘要 → 翻译 → 总结) | 各取所长,降低单次请求成本 |
|
||||
| **审批流水线** | 工作流中包含需要人工确认的步骤 | 兼顾自动化与安全合规 |
|
||||
| **数据管道** | ETL(抽取、转换、加载)流程的 Agent 化 | 用自然语言驱动数据处理 |
|
||||
|
||||
**性能优化:**
|
||||
### 性能优化建议
|
||||
|
||||
- 调整 MaxConcurrency 控制并发
|
||||
- 使用缓存避免重复计算
|
||||
- 流式输出提升用户体验
|
||||
1. **调整 MaxConcurrency**:根据 LLM API 的速率限制调参,一般 3~10 为宜
|
||||
2. **缓存层**:对相同 input + question 组合的结果做缓存,避免重复计算
|
||||
3. **自适应 chunk 大小**:根据文档类型(代码、散文、日志)动态调整分块策略
|
||||
4. **Early Exit**:当 top-1 分数远高于第二名时,跳过 filter 直接回答
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/_index]]
|
||||
- [[Eino/quick_start/chapter_04_tool_and_filesystem]] — 简单 Tool 的创建与文件系统访问(第二章的工具章节)
|
||||
- [[Eino/quick_start/chapter_07_interrupt_resume]] — Interrupt/Resume 机制(上一章)
|
||||
- [[Eino/quick_start/chapter_09_skill_console]] — Skill 系统(下一章)
|
||||
|
||||
@@ -1,25 +1,19 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-16"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第十章:A2UI 协议(流式 UI 组件)
|
||||
weight: 10
|
||||
tags: [eino, ai-development, go, quickstart, a2ui, sse]
|
||||
create time: 2026-04-29 16:00
|
||||
---
|
||||
|
||||
本章目标:实现 A2UI 协议,将 Agent 的输出渲染为流式 UI 组件。
|
||||
# 第十章:A2UI 协议(流式 UI 组件)(最终章)
|
||||
|
||||
## 重要说明:A2UI 的边界
|
||||
## 概述
|
||||
|
||||
A2UI 并不属于 Eino 框架本身的范畴,它是一个业务层的 UI 协议/渲染方案。本章把 A2UI 集成进前面章节逐步构建出来的 Agent,是为了提供一个端到端、可落地的完整示例:从模型调用、工具调用、工作流编排,到最终把结果以更友好的 UI 方式呈现出来。
|
||||
本章作为 ChatWithEino Quickstart 的最终章,引入 **A2UI 协议**——把 Agent 的事件流以 JSONL/SSE 的形式推送到前端,渲染为可增量更新的 UI 组件树。你将掌握 Agent 到 Web 的端到端集成方案,理解为什么 AI 应用需要从纯文本走向结构化、可交互的 UI 呈现。
|
||||
|
||||
在真实业务场景中,你完全可以根据产品形态选择不同的 UI 形式,例如:
|
||||
> [!tip] 一句话理解 A2UI
|
||||
>
|
||||
> **A2UI = Agent 输出 × UI 组件映射**。它定义了"Agent 做了什么"如何变成"用户看到了什么":文本 → Text 组件、工具调用 → Chip 卡片、进度更新 → 实时更新……一切通过声明式的组件树实现。
|
||||
|
||||
- Web / App:自定义组件、表格、卡片、图表等
|
||||
- IM/办公套件:消息卡片、交互式表单
|
||||
- 命令行:纯文本或 TUI(终端 UI)
|
||||
|
||||
Eino 更关注“可组合的智能执行与编排能力”,至于“如何呈现给用户”,属于业务层可以自由扩展的一环。
|
||||
---
|
||||
|
||||
## 代码位置
|
||||
|
||||
@@ -32,11 +26,11 @@ Eino 更关注“可组合的智能执行与编排能力”,至于“如何呈
|
||||
|
||||
## 前置条件
|
||||
|
||||
与第一章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)
|
||||
与第一章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)。
|
||||
|
||||
## 运行
|
||||
|
||||
在 `quickstart/chatwitheino` 目录下执行:
|
||||
在 `examples/quickstart/chatwitheino` 目录下执行:
|
||||
|
||||
```bash
|
||||
go run .
|
||||
@@ -48,50 +42,70 @@ go run .
|
||||
starting server on http://localhost:8080
|
||||
```
|
||||
|
||||
### (可选)启用 ch09 的 skills 能力
|
||||
启动后浏览器访问 `http://localhost:8080` 即可看到完整的 A2UI 交互界面。
|
||||
|
||||
最终 Web 版使用的 Agent 构建逻辑与 Chapter 9 对齐:当 `EINO_EXT_SKILLS_DIR` 指向一个合法 skills 目录时,会自动注册 `skill` 中间件,模型就能按需调用 `skill` 工具加载 `eino-guide` / `eino-component` / `eino-compose` / `eino-agent`。
|
||||
### (可选)启用 Skills 能力
|
||||
|
||||
最终 Web 版使用的 Agent 构建逻辑与第九章对齐:当 `EINO_EXT_SKILLS_DIR` 指向一个合法 skills 目录时,会自动注册 `skill` 中间件,模型就能按需调用 `skill` 工具加载文档。
|
||||
|
||||
```bash
|
||||
go run ./scripts/sync_eino_ext_skills.go -src /path/to/eino-ext -dest ./skills/eino-ext -clean
|
||||
EINO_EXT_SKILLS_DIR="$(pwd)/skills/eino-ext" go run .
|
||||
```
|
||||
|
||||
## 从文本到 UI:为什么需要 A2UI
|
||||
## A2UI 的定位与边界
|
||||
|
||||
前八章我们实现的 Agent 只输出文本,但现代 AI 应用需要更丰富的交互。
|
||||
> [!important] A2UI 不属于 Eino 框架本身
|
||||
|
||||
A2UI 是一个**业务层的 UI 协议/渲染方案**,不是 Eino 的核心 Component。本章把它集成进前面章节逐步构建出来的 Agent,是为了提供一个端到端、可落地的完整示例:从模型调用、工具调用、工作流编排,到最终把结果以更友好的 UI 方式呈现出来。
|
||||
|
||||
**真实业务场景中,你完全可以根据产品形态选择不同的 UI 形式:**
|
||||
|
||||
| 场景 | UI 形式 | 说明 |
|
||||
|------|---------|------|
|
||||
| Web / App | 自定义组件、表格、卡片、图表 | 最典型的 B/S 架构应用 |
|
||||
| IM / 办公套件 | 消息卡片、交互式表单 | 飞书、钉钉等平台的富消息 |
|
||||
| 命令行 | 纯文本或 TUI | Console 版 Agent 的原生形式 |
|
||||
|
||||
Eino 关注「可组合的智能执行与编排能力」,而「如何呈现给用户」属于业务层可以自由扩展的一环。
|
||||
|
||||
## 从纯文本到结构化的 UI:为什么需要 A2UI
|
||||
|
||||
> [!question] 思考一下
|
||||
>
|
||||
> 如果你要为一个 AI 聊天产品设计更丰富的交互体验,纯文本回复会遇到哪些瓶颈?
|
||||
|
||||
**纯文本输出的局限:**
|
||||
|
||||
- 无法展示结构化数据(表格、列表、卡片等)
|
||||
- 无法实时更新(进度条、状态变化等)
|
||||
- 无法嵌入交互元素(按钮、表单、链接等)
|
||||
- 无法支持多媒体(图片、视频、音频等)
|
||||
- ❌ 无法展示结构化数据(表格、列表、卡片等)
|
||||
- ❌ 无法实时更新(进度条、状态变化等)
|
||||
- ❌ 无法嵌入交互元素(按钮、表单、链接等)
|
||||
- ❌ 无法支持多媒体(图片、视频、音频等)
|
||||
|
||||
**A2UI 的定位:**
|
||||
|
||||
- **A2UI 是 Agent 到 UI 的协议**:定义了 Agent 输出如何映射到 UI 组件
|
||||
- **A2UI 支持流式渲染**:组件可以实时更新,无需等待完整响应
|
||||
- **A2UI 是声明式的**:Agent 只需声明"显示什么",UI 负责渲染
|
||||
- ✅ **协议映射**:Agent 输出 → UI 组件的声明式映射关系
|
||||
- ✅ **流式渲染**:组件实时更新,无需等待完整响应
|
||||
- ✅ **增量更新**:基于 dataKey 的数据绑定,文本流可逐 token 更新
|
||||
|
||||
**简单类比:**
|
||||
|
||||
- **纯文本输出** = "终端命令行"(只能显示文本)
|
||||
- **A2UI** = "Web 应用"(可以显示任何 UI 组件)
|
||||
|
||||
## 关键概念
|
||||
## A2UI v0.8 子集(本示例的边界)
|
||||
|
||||
### A2UI v0.8 子集(本示例的边界)
|
||||
|
||||
本 quickstart 并没有实现一个“完整的 A2UI 标准库”,而是实现了一个 **A2UI v0.8 的子集**:目标是把 Agent 的事件流,以稳定、可增量渲染的 UI 组件树方式推给浏览器。
|
||||
本 quickstart 并没有实现一个"完整的 A2UI 标准库",而是实现了一个 **A2UI v0.8 的子集**:目标是把 Agent 的事件流,以稳定、可增量渲染的 UI 组件树方式推给浏览器。
|
||||
|
||||
当前实现的 A2UI 消息类型与组件类型,以 [a2ui/types.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/a2ui/types.go) 为准。
|
||||
|
||||
### A2UI 消息:BeginRendering / SurfaceUpdate / DataModelUpdate / InterruptRequest
|
||||
### A2UI 消息类型:信封结构
|
||||
|
||||
每一行 SSE(`data: {...}`)承载一个 A2UI Message,Message 是一个“信封结构”,每次只会出现一个字段:
|
||||
每一行 SSE(`data: {...}`)承载一个 A2UI Message,Message 是一个"信封结构",每次只会出现一个字段:
|
||||
|
||||
**关键代码片段(注意:这是简化后的代码片段,不能直接运行,完整代码请参考 a2ui/types.go):**
|
||||
> [!note] 关键代码片段
|
||||
>
|
||||
> 注意:这是简化后的代码片段,不能直接运行,完整代码请参考 [a2ui/types.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/a2ui/types.go)。
|
||||
|
||||
```go
|
||||
type Message struct {
|
||||
@@ -103,52 +117,122 @@ type Message struct {
|
||||
}
|
||||
```
|
||||
|
||||
其中:
|
||||
| 消息类型 | 作用 | 触发时机 |
|
||||
|----------|------|----------|
|
||||
| `BeginRendering` | 告诉前端"开始渲染一个 surface",指定根节点 ID | 新会话开始时 |
|
||||
| `SurfaceUpdate` | 新增/更新一批组件(组件是树,用 id 互相引用) | 创建/修改 UI 结构时 |
|
||||
| `DataModelUpdate` | 更新 data bindings(用于流式文本增量渲染) | assistant 生成文本时 |
|
||||
| `InterruptRequest` | 通知前端展示批准/拒绝入口 | Agent 需要人类审批时 |
|
||||
| `DeleteSurface` | 删除某个 surface | 清理/重置会话时 |
|
||||
|
||||
- `BeginRendering`:告诉前端“开始渲染一个 surface(会话)”,并指定根节点 ID
|
||||
- `SurfaceUpdate`:新增/更新一批组件(组件是一个树,用 `id` 互相引用)
|
||||
- `DataModelUpdate`:更新 data bindings(用于把流式文本增量更新到某个 Text 组件)
|
||||
- `InterruptRequest`:当 Agent 触发 interrupt(例如审批)时,通知前端展示批准/拒绝入口
|
||||
|
||||
### A2UI 组件:Text / Column / Card / Row
|
||||
### A2UI 组件类型
|
||||
|
||||
本示例 UI 组件只实现了 4 种(见 [a2ui/types.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/a2ui/types.go)):
|
||||
|
||||
- `Text`:文本渲染(支持 `usageHint` 区分 caption/body/title);当 `dataKey` 存在时,文本来自 `DataModelUpdate`
|
||||
- `Column` / `Row`:布局(children 是组件 ID 列表)
|
||||
- `Card`:卡片容器(children 是组件 ID 列表)
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ComponentTree {
|
||||
<<abstract>>
|
||||
+id string
|
||||
+children []string
|
||||
}
|
||||
class TextComponent {
|
||||
+text string
|
||||
+dataKey string
|
||||
+usageHint string
|
||||
}
|
||||
class ColumnLayout {
|
||||
+spacing float64
|
||||
+align ItemsAlign
|
||||
}
|
||||
class RowLayout {
|
||||
+spacing float64
|
||||
+align ItemsAlign
|
||||
}
|
||||
class CardContainer {
|
||||
+title string
|
||||
+border bool
|
||||
}
|
||||
ComponentTree <|-- TextComponent
|
||||
ComponentTree <|-- ColumnLayout
|
||||
ComponentTree <|-- RowLayout
|
||||
ComponentTree <|-- CardContainer
|
||||
note for TextComponent "支持 dataKey\n流式绑定"
|
||||
note for ColumnLayout "垂直布局容器"
|
||||
note for RowLayout "水平布局容器"
|
||||
note for CardContainer "内容容器\n无布局功能"
|
||||
```
|
||||
|
||||
## A2UI 的实现:把 AgentEvent 转成 A2UI SSE
|
||||
各组件职责:
|
||||
|
||||
最终 Web 版的核心链路是:
|
||||
| 组件 | 用途 | 特性 |
|
||||
|------|------|------|
|
||||
| `Text` | 文本渲染 | 支持 `usageHint`(caption/body/title),当存在 `dataKey` 时文本来自 `DataModelUpdate` |
|
||||
| `Column` | 垂直布局 | children 是组件 ID 列表 |
|
||||
| `Row` | 水平布局 | children 是组件 ID 列表 |
|
||||
| `Card` | 卡片容器 | children 是组件 ID 列表,仅做视觉分组 |
|
||||
|
||||
- 后端运行 Agent,得到 `*adk.AsyncIterator[*adk.AgentEvent]`
|
||||
- 把事件流转换为 A2UI JSONL/SSE 流输出给浏览器(见 [a2ui/streamer.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/a2ui/streamer.go))
|
||||
- 前端解析 SSE 的 `data:` 行并渲染组件树(见 [static/index.html](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/static/index.html))
|
||||
> [!warning] Card ≠ 布局容器
|
||||
>
|
||||
> `Card` 不提供任何布局控制(不排布子元素的位置),它只是一个有视觉边界的容器。如需布局请用 `Column` / `Row`。
|
||||
|
||||
## A2UI 的实现链路
|
||||
|
||||
最终 Web 版的核心链路是三阶段管道:
|
||||
|
||||
1. **后端运行 Agent**:得到 `*adk.AsyncIterator[*adk.AgentEvent]`
|
||||
2. **事件 → A2UI JSONL/SSE 流**:转换为 A2UI 消息推送给浏览器(见 [a2ui/streamer.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/a2ui/streamer.go))
|
||||
3. **前端解析并渲染**:读取 SSE 流并渲染组件树(见 [static/index.html](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/static/index.html))
|
||||
|
||||
### 服务端路由(高层)
|
||||
|
||||
与 A2UI 相关的关键接口(见 [server/server.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/server/server.go)):
|
||||
|
||||
- `GET /`:返回前端页面 `static/index.html`
|
||||
- `POST /sessions/:id/chat`:返回 SSE 流(A2UI messages),把 Agent 运行结果边跑边渲染到 UI
|
||||
- `GET /sessions/:id/render`:返回 JSONL(A2UI messages),用于“选中会话时回放历史”
|
||||
- `POST /sessions/:id/approve`:处理 interrupt 的批准/拒绝并继续返回 SSE 流
|
||||
| 方法 | 路径 | 响应 | 说明 |
|
||||
|------|------|------|------|
|
||||
| GET | `/` | HTML | 返回前端页面 |
|
||||
| POST | `/sessions/:id/chat` | SSE 流 | Agent 运行结果实时渲染 |
|
||||
| GET | `/sessions/:id/render` | JSONL | 回放历史消息 |
|
||||
| POST | `/sessions/:id/approve` | SSE 流 | interrupt 批准后继续执行 |
|
||||
|
||||
### 事件流转换(高层)
|
||||
### 事件流转换
|
||||
|
||||
服务端把 `Runner.Run(...)` 的事件流交给 `a2ui.StreamToWriter(...)`,后者负责:
|
||||
|
||||
- 对 user/assistant/tool 的输出做拆分
|
||||
- 把 tool call / tool result 渲染成 “chip 卡片”
|
||||
- 把 assistant 的流式 token 做成 `DataModelUpdate`,实现“边生成边渲染”
|
||||
- 遇到 interrupt 时发送 `InterruptRequest`,并暂停等待人类批准
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["AgentEvent 输入"] --> B{事件类型}
|
||||
|
||||
## 前端集成:fetch + SSE(不是 WebSocket)
|
||||
B -->|"user"| C["渲染 User 气泡"]
|
||||
B -->|"assistant"| D["创建 DataModelUpdate\n流式追加文本"]
|
||||
B -->|"tool call"| E["渲染 ToolCall Chip 卡片"]
|
||||
B -->|"tool result"| F["渲染 ToolResult Chip 卡片"]
|
||||
B -->|"interrupt"| G["发送 InterruptRequest\n暂停等待人类审批"]
|
||||
|
||||
- 前端通过 `fetch('/sessions/:id/chat')` 发起请求,然后从 `res.body` 读取流式字节,按行切分并解析 `data: {...}` 的 JSON(见 [static/index.html](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/static/index.html))。
|
||||
C --> H["SSE JSONL 输出"]
|
||||
D --> H
|
||||
E --> H
|
||||
F --> H
|
||||
G --> H
|
||||
|
||||
**关键代码片段(注意:这是简化后的代码片段,不能直接运行,完整代码请参考 static/index.html):**
|
||||
style D fill:#fff3e0
|
||||
style G fill:#fce4ec
|
||||
```
|
||||
|
||||
核心处理逻辑:
|
||||
|
||||
- **User 输出**:渲染为用户消息气泡
|
||||
- **Assistant 流式 token**:创建 `DataModelUpdate`,通过 `dataKey` 绑定到一个 `Text` 组件上,实现"边生成边渲染"
|
||||
- **Tool Call / Tool Result**:渲染为独立的 chip 卡片
|
||||
- **Interrupt**:发送 `InterruptRequest`,暂停等待人类批准后 resume
|
||||
|
||||
### 前端集成:Fetch + SSE(不是 WebSocket)
|
||||
|
||||
前端通过 `fetch('/sessions/:id/chat')` 发起请求,然后从 `res.body` 读取流式字节,按行切分并解析 `data: {...}` 的 JSON:
|
||||
|
||||
> [!note] 前端代码片段
|
||||
>
|
||||
> 注意:这是简化后的代码片段,不能直接运行,完整代码请参考 [static/index.html](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/static/index.html)。
|
||||
|
||||
```javascript
|
||||
const res = await fetch(`/sessions/${id}/chat`, {
|
||||
@@ -176,77 +260,107 @@ while (true) {
|
||||
}
|
||||
```
|
||||
|
||||
## A2UI 流式渲染流程(概览)
|
||||
> [!tip] 为什么用 SSE 而不是 WebSocket?
|
||||
>
|
||||
> - **SSE**(Server-Sent Events)是单向的、HTTP 兼容、天然支持断线重连语义
|
||||
> - 对于"服务器推送 UI 事件,客户端只需消费"的场景,SSE 比 WebSocket 更轻量
|
||||
> - 如果未来需要双向交互(如键盘快捷键、实时光标同步),才考虑升级到 WebSocket
|
||||
|
||||
## A2UI 流式渲染流程
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant U as 用户
|
||||
participant FE as 前端浏览器
|
||||
participant BE as 后端 Server
|
||||
participant AG as Agent
|
||||
participant LM as LLM
|
||||
|
||||
U->>FE: 输入消息
|
||||
FE->>BE: POST /sessions/:id/chat
|
||||
BE->>AG: Runner.Run()
|
||||
AG->>LM: 发送请求
|
||||
LM-->>AG: token 流式返回
|
||||
loop 每个 AgentEvent
|
||||
AG-->>BE: AgentEvent
|
||||
alt assistant token
|
||||
BE->>FE: DataModelUpdate (dataKey 绑定)
|
||||
else tool call
|
||||
BE->>FE: SurfaceUpdate (Chip 卡片)
|
||||
end
|
||||
FE-->>U: UI 增量更新
|
||||
end
|
||||
AG-->>AG: 可能需要 Interrupt
|
||||
AG->>BE: InterruptRequest
|
||||
BE->>FE: InterruptRequest
|
||||
FE-->>U: 展示审批按钮
|
||||
U->>FE: 点击批准
|
||||
FE->>BE: POST /sessions/:id/approve
|
||||
BE->>AG: Resume 继续执行
|
||||
```
|
||||
┌─────────────────────────────────────────┐
|
||||
│ 用户:分析这个文件 │
|
||||
└─────────────────────────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Agent 开始处理 │
|
||||
│ A2UI: AddText │
|
||||
│ "正在分析..." │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 调用 Tool │
|
||||
│ A2UI: AddProgress │
|
||||
│ 进度: 0% │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Tool 执行中 │
|
||||
│ A2UI: UpdateProgress│
|
||||
│ 进度: 50% │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ Tool 完成 │
|
||||
│ A2UI: tool result │
|
||||
└──────────────────────┘
|
||||
↓
|
||||
┌──────────────────────┐
|
||||
│ 显示结果 │
|
||||
│ A2UI: DataModelUpdate│
|
||||
│ (流式更新 assistant)│
|
||||
└──────────────────────┘
|
||||
```
|
||||
|
||||
## 从 Quickstart 到生产落地
|
||||
|
||||
> [!tip] 可扩展的设计思路
|
||||
>
|
||||
> A2UI 的协议层和 Agent 层是解耦的。你可以只做其中的任意一部分。
|
||||
|
||||
| 方向 | 替换方案 | 适用场景 |
|
||||
|------|---------|----------|
|
||||
| 不同 UI 形态 | React/Vue 组件、移动端原生组件、TUI | 产品定位差异 |
|
||||
| 传输协议升级 | gRPC Streaming / GraphQL Subscription | 需要更强的双向交互 |
|
||||
| 渲染引擎切换 | 服务端 SSR → 客户端动态渲染 | CDN 加速需求 |
|
||||
| 多模态扩展 | 嵌入图片、图表、代码高亮 | 数据类/分析型 Agent |
|
||||
|
||||
## 本章小结
|
||||
|
||||
- **A2UI**:Agent 到 UI 的协议,定义了 Agent 输出如何映射到 UI 组件
|
||||
- **子集实现**:本示例只实现了 Text/Column/Card/Row 与 data binding
|
||||
- **流式输出**:后端以 SSE 推送 A2UI JSONL,前端增量渲染组件树
|
||||
- **事件到 UI**:把 `AgentEvent` 转为 `tool call / tool result / assistant stream` 的可视化输出
|
||||
| 核心概念 | 说明 | 关键点 |
|
||||
|----------|------|--------|
|
||||
| **A2UI 协议** | Agent 到 UI 的映射协议 | 声明式组件树 + 数据绑定 |
|
||||
| **消息信封** | BeginRendering / SurfaceUpdate / DataModelUpdate | 每种消息对应不同生命周期 |
|
||||
| **组件系统** | Text / Column / Card / Row | 4 种基础组件覆盖常见 UI 场景 |
|
||||
| **SSE 流** | AgentEvent → A2UI JSONL → 前端渲染 | 单向推送、增量更新 |
|
||||
| **中断协作** | InterruptRequest + approve 机制 | 人机协同的关键路径 |
|
||||
|
||||
## 系列收尾:这个 Quickstart Agent 的完整愿景
|
||||
> [!success] 学习成果
|
||||
>
|
||||
> 完成本章后,你应该能够:
|
||||
> - 理解 A2UI 协议的设计思路和适用场景
|
||||
> - 掌握 Agent 事件流到前端 UI 的完整转换链路
|
||||
> - 使用 SSE 实现增量渲染的前端集成方案
|
||||
> - 知道如何将这个骨架扩展到不同的产品形态中
|
||||
|
||||
到本章为止,我们用一个可以实际运行的 Agent 串起了 Eino 的核心能力。你可以把它理解为一个可扩展的“端到端 Agent 应用骨架”:
|
||||
## 下一章预告
|
||||
|
||||
- 运行时:Runner 驱动执行,支持流式输出与事件模型
|
||||
- 工具层:Filesystem / Shell 等 Tool 能力接入,工具错误可被安全处理
|
||||
- 中间件:可插拔的 middleware/handler,用于错误处理、重试、审批等横切能力
|
||||
- 可观测:callbacks/trace 能力把关键链路打通,便于调试与线上观测
|
||||
- 人机协作:interrupt/resume + checkpoint 支持审批、补参、分支选择等交互式流程
|
||||
- 确定性编排:compose(graph/chain/workflow)把复杂业务流程组织为可维护、可复用的执行图
|
||||
- 业务交付:像 A2UI 这样的 UI 集成,属于业务层自由选择的一环,用来把 Agent 能力以合适的产品形态呈现给用户
|
||||
这是 Quickstart 系列的终章。后续如果你想深入:
|
||||
|
||||
你可以在这个骨架上逐步替换/扩展任意环节:模型、工具、存储、工作流、前端渲染协议,而不需要推倒重来。
|
||||
- [[Eino/quick_start/chapter_09_skill_console]] — 回顾第九章的 Skill 知识注入能力
|
||||
- [[Eino/README]] — 探索 Eino 框架的系统性学习路径
|
||||
|
||||
## 扩展思考
|
||||
|
||||
**其他组件类型:**
|
||||
### 其他组件类型(可选实现方向)
|
||||
|
||||
- 图表组件(折线图、柱状图、饼图)
|
||||
- 地图组件
|
||||
- 时间线组件
|
||||
- 树形组件
|
||||
- 标签页组件
|
||||
| 组件 | 说明 | 实现难度 |
|
||||
|------|------|----------|
|
||||
| **图表组件** | 折线图、柱状图、饼图 | 中高(需接入图表库) |
|
||||
| **地图组件** | 地理信息可视化 | 中(依赖地图 SDK) |
|
||||
| **时间线组件** | 事件顺序排列展示 | 低(已有 Column 即可) |
|
||||
| **树形组件** | 层级数据结构展示 | 中(递归渲染逻辑) |
|
||||
| **标签页组件** | 多面板 Tab 切换 | 低(State 管理即可) |
|
||||
|
||||
**高级功能:**
|
||||
### 高级交互能力
|
||||
|
||||
- 组件交互(点击、拖拽、输入)
|
||||
- 条件渲染
|
||||
- 组件动画
|
||||
- 响应式布局
|
||||
| 能力 | 说明 | 技术要点 |
|
||||
|------|------|----------|
|
||||
| **组件交互** | 点击、拖拽、输入反馈 | 前端事件 → API 调用 |
|
||||
| **条件渲染** | 根据数据决定组件显隐 | 服务端根据状态发送不同 SurfaceUpdate |
|
||||
| **组件动画** | 平滑过渡效果 | CSS transition / animation |
|
||||
| **响应式布局** | 自适应屏幕尺寸 | 媒体查询 + 弹性布局 |
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/_index]] — Quickstart 系列统一入口
|
||||
- [[Eino/quick_start/chapter_08_graph_tool]] — Graph Tool 复杂工作流编排
|
||||
- [[Eino/quick_start/chapter_09_skill_console]] — Skill 知识与指令注入
|
||||
- [[Eino/quick_start/chapter_07_interrupt_resume]] — Interrupt 与 Resume 机制
|
||||
|
||||
@@ -1,61 +1,76 @@
|
||||
---
|
||||
Description: ""
|
||||
date: "2026-03-24"
|
||||
lastmod: ""
|
||||
tags: []
|
||||
title: 第九章:Skill(Console)
|
||||
weight: 9
|
||||
create time: 2026-04-29 15:30
|
||||
---
|
||||
|
||||
本章目标:在第八章(RAG + Interrupt/Resume + Checkpoint)基础上,引入 `skill` 中间件,让 Agent 可以发现并加载一组可复用的技能文档(`SKILL.md`),并在需要时通过工具调用使用它们。
|
||||
# 第九章:Skill(Console)
|
||||
|
||||
## 代码位置
|
||||
## 概述
|
||||
|
||||
- 入口代码:[cmd/ch09/main.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/cmd/ch09/main.go)
|
||||
- 同步脚本:[scripts/sync_eino_ext_skills.go](https://github.com/cloudwego/eino-examples/blob/main/quickstart/chatwitheino/scripts/sync_eino_ext_skills.go)
|
||||
本章在上一章(RAG + Interrupt/Resume + Checkpoint)的基础上,引入 **Skill** 中间件。通过 Skill 机制,Agent 可以发现并加载一组可复用的"技能文档"(`SKILL.md`),并在需要时自动调用它们——让 Agent 获得结构化的领域知识,而不需要把所有知识写进系统提示词里。
|
||||
|
||||
> [!TIP] 核心目标
|
||||
> 学会用 Skill 中间件把一个稳定的知识集合注入到 Agent 中,并理解 Skill 与 Tool 的区别、注册方式、以及验证方法。
|
||||
|
||||
---
|
||||
|
||||
## 前置条件
|
||||
|
||||
- 与第一章一致:需要配置一个可用的 ChatModel(OpenAI 或 Ark)
|
||||
- 准备好 `eino-ext` PR 提供的 skills(`eino-guide` / `eino-component` / `eino-compose` / `eino-agent`)
|
||||
- 准备好 `eino-ext` PR 提供的 skills 资源:`eino-guide` / `eino-component` / `eino-compose` / `eino-agent`
|
||||
|
||||
为什么是这四个?
|
||||
> [!QUESTION] 为什么是这四个 skill?
|
||||
>
|
||||
> ChatWithEino 的定位是「帮用户学习 Eino 框架、并尝试用 AI 辅助写 Eino 代码」。这四个 skill 恰好覆盖了关键知识点:
|
||||
>
|
||||
> - **`eino-guide`** — 学习入口与导航(从哪里开始、怎么快速跑起来)
|
||||
> - **`eino-component`** — Component 接口与各类实现参考(Model / Embedding / Retriever / Tool / Callback 等)
|
||||
> - **`eino-compose`** — 编排与确定性工作流参考(Graph / Chain / Workflow 等)
|
||||
> - **`eino-agent`** — ADK / Agent 相关参考(Agent / Runner / Middleware / Filesystem / Human-in-the-loop 等)
|
||||
|
||||
ChatWithEino 的定位是“帮用户学习 Eino 框架、并尝试用 AI 辅助写 Eino 代码”。这四个 skills 正好覆盖了这个目标所需的关键知识面:
|
||||
Skills 来源可以是:
|
||||
|
||||
- `eino-guide`:学习入口与导航(从哪里开始、怎么快速跑起来)
|
||||
- `eino-component`:Component 接口与各类实现参考(Model/Embedding/Retriever/Tool/Callback 等)
|
||||
- `eino-compose`:编排与确定性工作流参考(Graph/Chain/Workflow 等)
|
||||
- `eino-agent`:ADK/Agent 相关参考(Agent、Runner、Middleware、Filesystem、Human-in-the-loop 等)
|
||||
- `eino-ext` 仓库本地路径(同步脚本会自动读取 `<src>/skills/...`)
|
||||
- 你已安装 skills 的目录(目录下能看到上述四个子目录)
|
||||
|
||||
skills 的来源可以是:
|
||||
---
|
||||
|
||||
- `eino-ext` 仓库本地路径(脚本会自动读取 `<src>/skills/...`)
|
||||
- 或你已安装 skills 的目录(目录下能看到上述四个子目录)
|
||||
## 正文
|
||||
|
||||
## 从 Graph Tool 到 Skill:为什么需要“技能文档”
|
||||
### 从 Graph Tool 到 Skill:为什么需要"技能文档"
|
||||
|
||||
第八章解决的是“复杂工作流如何做成一个可调用的 Tool”(Graph Tool)。但你在构建一个面向框架学习/开发辅助的 Agent 时,还会遇到另一类问题:**如何把一组稳定、可复用的知识与指令注入到 Agent 里,并让它在运行时按需加载?**
|
||||
第八章我们解决了「复杂工作流如何做成一个可调用的 Tool」的问题。但当你构建一个面向框架学习/开发辅助的 Agent 时,还会遇到另一类挑战:
|
||||
|
||||
这就是 Skill 的定位:
|
||||
> **如何把一组稳定、可复用的知识与指令注入到 Agent 里,并让它在运行时按需加载?**
|
||||
|
||||
- **Tool** 更像“动作/能力”:读文件、跑 workflow、调用外部系统
|
||||
- **Skill** 更像“可复用的知识/指令包”:用一组 markdown(`SKILL.md` + `reference/*.md`)描述“如何做某类事”
|
||||
这就是 Skill 的切入点:
|
||||
|
||||
简单类比:
|
||||
- **Tool** = "能做什么"(函数/接口级别的能力)
|
||||
- **Skill** = "怎么做"(可复用的说明书/操作手册)
|
||||
|
||||
- **Tool** = “能做什么”(函数/接口)
|
||||
- **Skill** = “怎么做”(可复用的说明书/操作手册)
|
||||
```mermaid
|
||||
graph LR
|
||||
A["Agent"] --> B["Tool 层<br/>读文件 / 执行流程 / 调外部API"]
|
||||
A --> C["Skill 层<br/>知识文档 / 最佳实践 / 操作手册"]
|
||||
C --> D["eino-guide<br/>学习入口"]
|
||||
C --> E["eino-component<br/>组件参考"]
|
||||
C --> F["eino-compose<br/>编排参考"]
|
||||
C --> G["eino-agent<br/>ADK参考"]
|
||||
```
|
||||
|
||||
## 运行
|
||||
简单说:**Skill 是一种可被模型发现的结构化知识包**。每个 Skill 以 `SKILL.md` 为核心描述文件,辅以 `reference/*.md` 参考资料。
|
||||
|
||||
在 `quickstart/chatwitheino` 目录下执行:
|
||||
### 运行步骤
|
||||
|
||||
### 1) 同步 eino-ext skills 到本地目录
|
||||
在 `quickstart/chatwitheino` 目录下执行以下两步:
|
||||
|
||||
为了让 `skill` 中间件可以“发现”这些 skills,需要把它们放到一个统一目录下,并满足扫描约定:
|
||||
#### 1) 同步 eino-ext skills 到本地目录
|
||||
|
||||
- `EINO_EXT_SKILLS_DIR/<skillName>/SKILL.md`
|
||||
为了让 `skill` 中间件可以"发现"这些 skills,需要把它们放到一个统一目录下,满足扫描约定:
|
||||
|
||||
```
|
||||
EINO_EXT_SKILLS_DIR/<skillName>/SKILL.md
|
||||
```
|
||||
|
||||
同步命令(推荐):
|
||||
|
||||
@@ -63,80 +78,120 @@ skills 的来源可以是:
|
||||
go run ./scripts/sync_eino_ext_skills.go -src /path/to/eino-ext -dest ./skills/eino-ext -clean
|
||||
```
|
||||
|
||||
说明:
|
||||
> [!NOTE] `-src` 参数说明
|
||||
>
|
||||
> - 形式一:`eino-ext` 仓库根目录 → 脚本自动读取 `<src>/skills/...`
|
||||
> - 形式二:你已安装 skills 的目录 → 要求目录下包含 `eino-guide/`、`eino-component/` 等子目录
|
||||
|
||||
- `-src` 支持两种形式:
|
||||
- `eino-ext` 仓库根目录(脚本会自动读取 `<src>/skills/...`)
|
||||
- 你已安装 skills 的目录(目录下应包含 `eino-guide/`、`eino-component/` 等子目录)
|
||||
- `-dest` 默认是 `./skills/eino-ext`(可以省略)
|
||||
|
||||
### 2) 启动 Chapter 9
|
||||
#### 2) 启动 Chapter 9
|
||||
|
||||
```bash
|
||||
EINO_EXT_SKILLS_DIR=/absolute/path/to/chatwitheino/skills/eino-ext go run ./cmd/ch09
|
||||
```
|
||||
|
||||
输出示例(节选):
|
||||
控制台输出示例:
|
||||
|
||||
```
|
||||
Skills dir: /.../skills/eino-ext
|
||||
Enter your message (empty line to exit):
|
||||
```
|
||||
|
||||
## 在 DeepAgent 中启用 Skill
|
||||
### 在 DeepAgent 中启用 Skill
|
||||
|
||||
本章的 “Skill 可被调用” 不是自动发生的,你需要在 Agent 构建时把 `skill` 中间件注册进去。核心就是三步:
|
||||
Skill 不会被自动加载 —— 你需要在 Agent 构建时显式注册 `skill` 中间件。核心三步:
|
||||
|
||||
1. 用本地 filesystem backend(本章用 `eino-ext/adk/backend/local`)提供文件读取/Glob 能力
|
||||
2. 用 `skill.NewBackendFromFilesystem` 把 `EINO_EXT_SKILLS_DIR` 变成一个 Skill Backend
|
||||
3. 用 `skill.NewMiddleware` 生成中间件,并把它塞进 DeepAgent 的 `Handlers`
|
||||
| 步骤 | 操作 | 关键 API |
|
||||
|------|------|----------|
|
||||
| 1️⃣ | 创建文件系统 backend | `localbk.NewBackend(ctx, &localbk.Config{})` |
|
||||
| 2️⃣ | 构建 Skill Backend | `skill.NewBackendFromFilesystem(ctx, cfg)` |
|
||||
| 3️⃣ | 生成中间件并注入 DeepAgent | `skill.NewMiddleware(ctx, cfg)` |
|
||||
|
||||
**关键代码片段(注意:这是简化后的代码片段,不能直接运行,完整代码请参考 ****cmd/ch09/main.go****):**
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant User as 用户
|
||||
participant Agent as DeepAgent
|
||||
participant SkillMW as Skill 中间件
|
||||
participant Backend as Skill Backend
|
||||
participant FS as 本地文件系统
|
||||
|
||||
User->>Agent: 发送消息
|
||||
Agent->>SkillMW: 处理请求
|
||||
SkillMW->>Backend: 按 skillName 查找 SKILL.md
|
||||
Backend->>FS: Glob / Read 文件
|
||||
FS-->>Backend: 返回 Markdown 内容
|
||||
Backend-->>SkillMW: 技能上下文
|
||||
SkillMW-->>Agent: 注入知识到 prompt
|
||||
Agent->>User: 返回回复
|
||||
```
|
||||
|
||||
**关键代码片段(简化版,完整代码见 `cmd/ch09/main.go`):**
|
||||
|
||||
```go
|
||||
// Step 1: 本地 filesystem backend
|
||||
backend, _ := localbk.NewBackend(ctx, &localbk.Config{})
|
||||
|
||||
// Step 2: 把 $EINO_EXT_SKILLS_DIR 变成 Skill Backend
|
||||
skillBackend, _ := skill.NewBackendFromFilesystem(ctx, &skill.BackendFromFilesystemConfig{
|
||||
Backend: backend,
|
||||
BaseDir: skillsDir, // = $EINO_EXT_SKILLS_DIR
|
||||
BaseDir: skillsDir, // = os.Getenv("EINO_EXT_SKILLS_DIR")
|
||||
})
|
||||
|
||||
// Step 3: 创建中间件并注册到 DeepAgent
|
||||
skillMiddleware, _ := skill.NewMiddleware(ctx, &skill.Config{
|
||||
Backend: skillBackend,
|
||||
})
|
||||
|
||||
agent, _ := deep.New(ctx, &deep.Config{
|
||||
ChatModel: cm,
|
||||
Backend: backend,
|
||||
Backend: backend,
|
||||
StreamingShell: backend,
|
||||
Handlers: []adk.ChatModelAgentMiddleware{
|
||||
skillMiddleware,
|
||||
// ... 其他中间件,比如 approval/safeTool/retry 等
|
||||
// ... 其他中间件(approval / safeTool / retry 等)
|
||||
},
|
||||
})
|
||||
```
|
||||
|
||||
补充说明:
|
||||
> [!WARNING] 容错设计
|
||||
>
|
||||
> 本 quickstart 保证了"没配置 skills 也能跑":代码中对 `EINO_EXT_SKILLS_DIR` 做了存在性检查,目录不存在则跳过注册 `skillMiddleware`。此时仍可正常对话和使用 RAG 工具。
|
||||
|
||||
- 本 quickstart 为了保证 “没配置 skills 也能跑”,在代码里对 `EINO_EXT_SKILLS_DIR` 做了存在性检查:目录存在才注册 `skillMiddleware`;否则跳过(此时仍可对话与使用 RAG 工具)。
|
||||
- Skill 工具的入参是一个 JSON:`{"skill": "<skillName>"}`,例如 `{"skill":"eino-guide"}`。
|
||||
### Skill 工具的入参格式
|
||||
|
||||
## 快速验证(推荐)
|
||||
Skill 被注册为 Tool 后,模型的调用入参是一个 JSON 对象:
|
||||
|
||||
启动后输入一条指令,明确要求模型调用 skill 工具(用于验证 skills 已被发现且可被加载):
|
||||
```json
|
||||
{"skill": "eino-guide"}
|
||||
```
|
||||
|
||||
其中 `"skill"` 键对应要激活的技能名称。
|
||||
|
||||
### 快速验证
|
||||
|
||||
启动后输入一条明确要求模型调用 skill 工具的指令:
|
||||
|
||||
```
|
||||
Use the skill tool with skill="eino-guide" and tell me what the entry point is for getting started.
|
||||
```
|
||||
|
||||
你应当能在控制台看到类似输出:
|
||||
你应该看到:
|
||||
|
||||
- `[tool result] Launching skill: eino-guide`
|
||||
- Tool result 中包含 `Base directory for this skill: .../eino-guide`
|
||||
- `[tool call] ...` — 模型发起了 skill 工具调用
|
||||
- `[tool result] Launching skill: eino-guide` — 技能被成功激活
|
||||
- Tool result 中包含 `Base directory for this skill: .../eino-guide` — 确认文件读取正确
|
||||
|
||||
## 你会看到什么
|
||||
### 会话恢复
|
||||
|
||||
- 当模型调用 skill 工具时,控制台会打印:
|
||||
- `[tool call] ...`
|
||||
- `[tool result] ...`(对结果做了截断展示)
|
||||
- 会话保存在 `SESSION_DIR`(默认 `./data/sessions`),支持恢复:
|
||||
- `go run ./cmd/ch09 --session <id>`
|
||||
会话数据保存在 `SESSION_DIR`(默认 `./data/sessions`),支持通过 `--session` 参数恢复:
|
||||
|
||||
```bash
|
||||
go run ./cmd/ch09 --session <session-id>
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[Eino/quick_start/chapter_08_graph_tool]] — 上一章:Graph Tool,理解 Tool 作为"动作能力"的基础
|
||||
- [[Eino/quick_start/chapter_05_middleware]] — Middleware 机制,所有中间件的通用注册方式
|
||||
- [[Eino/quick_start/chapter_04_tool_and_filesystem]] — Tool 与 Filesystem,文件系统 backend 的来源
|
||||
|
||||
Reference in New Issue
Block a user