524 lines
17 KiB
Markdown
524 lines
17 KiB
Markdown
|
|
# CamTalk Eino 框架与编排设计
|
|||
|
|
|
|||
|
|
## 1. 概述
|
|||
|
|
|
|||
|
|
### 1.1 为什么选择 Eino
|
|||
|
|
|
|||
|
|
[CloudWeGo Eino](https://github.com/cloudwego/eino) 是字节跳动 CloudWeGo 团队开源的 AI 应用开发框架,提供基于图(Graph)的编排能力、组件抽象和流式处理支持。
|
|||
|
|
|
|||
|
|
CamTalk 使用 Eino 替代原有的手写 goroutine 管道,实现 STT → LLM → TTS 的声明式编排。
|
|||
|
|
|
|||
|
|
**技术选型对比:**
|
|||
|
|
|
|||
|
|
| 维度 | 手写 goroutine(旧方案) | Eino Graph(新方案) |
|
|||
|
|
|------|------------------------|---------------------|
|
|||
|
|
| 编排方式 | 手动 `go func()` + `sync.WaitGroup` | 声明式 DAG,类型安全 |
|
|||
|
|
| 流式处理 | 自定义 `chan` 传递 | `StreamReader` + `Pipe`,自动转换 |
|
|||
|
|
| 错误处理 | 各节点独立处理,不一致 | Graph 级别统一错误传播 |
|
|||
|
|
| 回调/AOP | 日志散落各处 | `callbacks.Handler` 统一注入 |
|
|||
|
|
| 配置灵活性 | Pipeline 创建时固定 | 每请求 `Option` 动态注入 |
|
|||
|
|
| 可测试性 | 需启动 goroutine | `Graph.Invoke()` 直接测试 |
|
|||
|
|
| 扩展性 | 修改 Pipeline 代码 | 添加节点 + 边,无侵入 |
|
|||
|
|
| 并发安全 | 手动 `sync` | State 自动加锁 |
|
|||
|
|
|
|||
|
|
**选择 Eino 的核心理由:**
|
|||
|
|
1. Go 原生,泛型支持,编译时类型检查
|
|||
|
|
2. 原生流式处理(`StreamReader`),适合 LLM token 级推送
|
|||
|
|
3. Graph 支持分支、并行、循环,满足当前和未来需求
|
|||
|
|
4. Callback 机制实现 AOP(日志、指标、消息推送)
|
|||
|
|
5. eino-ext 提供 OpenAI ChatModel 实现,直接对接 DashScope
|
|||
|
|
|
|||
|
|
### 1.2 旧方案的问题
|
|||
|
|
|
|||
|
|
当前后端 AI 编排层(`internal/orchestrator/pipeline.go`)为手写 goroutine 管道存在以下问题:
|
|||
|
|
|
|||
|
|
1. **编排逻辑硬编码**:STT→LLM→TTS 流程写死,扩展困难
|
|||
|
|
2. **并发控制粗糙**:手动 goroutine 调度,缺乏结构化流式传递
|
|||
|
|
3. **无回调/AOP 机制**:日志、指标、追踪散落各处
|
|||
|
|
4. **配置耦合**:模型名、TTS 参数等硬编码在结构体
|
|||
|
|
5. **错误处理不一致**:TTS 错误静默吞掉,STT/LLM 错误通过 Sender 发送
|
|||
|
|
|
|||
|
|
### 1.3 核心依赖版本
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
github.com/cloudwego/eino v0.9.9
|
|||
|
|
github.com/cloudwego/eino-ext/components/model/openai v0.1.13
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
## 2. Eino 核心概念
|
|||
|
|
|
|||
|
|
### 2.1 Lambda
|
|||
|
|
|
|||
|
|
Lambda 是 Graph 中的可组合函数单元,支持四种模式:
|
|||
|
|
|
|||
|
|
| 模式 | 函数签名 | 构造方法 | 说明 |
|
|||
|
|
|------|---------|---------|------|
|
|||
|
|
| Invoke | `I → O` | `compose.InvokableLambda()` | 同步调用 |
|
|||
|
|
| Stream | `I → StreamReader[O]` | `compose.StreamableLambda()` | 流式输出 |
|
|||
|
|
| Collect | `StreamReader[I] → O` | `compose.CollectableLambda()` | 流式输入 |
|
|||
|
|
| Transform | `StreamReader[I] → StreamReader[O]` | `compose.TransformableLambda()` | 双向流式 |
|
|||
|
|
|
|||
|
|
**返回类型**:所有 Lambda 构造函数返回 `*compose.Lambda`。
|
|||
|
|
|
|||
|
|
### 2.2 Graph
|
|||
|
|
|
|||
|
|
Graph 是有向无环图(DAG)编排器,支持:
|
|||
|
|
- **节点**:Lambda、ChatModel、ToolsNode 等
|
|||
|
|
- **边**:`g.AddEdge(from, to)` 定义数据流向
|
|||
|
|
- **分支**:`g.AddBranch()` 条件路由
|
|||
|
|
- **State**:`compose.WithGenLocalState()` 跨节点共享状态
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
g := compose.NewGraph[PipelineInput, PipelineOutput]()
|
|||
|
|
g.AddLambdaNode("stt", sttLambda)
|
|||
|
|
g.AddChatModelNode("llm", chatModel)
|
|||
|
|
g.AddEdge(compose.START, "stt")
|
|||
|
|
g.AddEdge("stt", "llm")
|
|||
|
|
g.AddEdge("llm", compose.END)
|
|||
|
|
|
|||
|
|
runnable, err := g.Compile(ctx)
|
|||
|
|
output, err := runnable.Invoke(ctx, input) // 同步调用
|
|||
|
|
stream, err := runnable.Stream(ctx, input) // 流式调用
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 2.3 ChatModel
|
|||
|
|
|
|||
|
|
ChatModel 是 LLM 组件抽象,接口定义:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
type BaseChatModel interface {
|
|||
|
|
Generate(ctx, []*schema.Message, ...Option) (*schema.Message, error)
|
|||
|
|
Stream(ctx, []*schema.Message, ...Option) (*schema.StreamReader[*schema.Message], error)
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
CamTalk 使用 `eino-ext/components/model/openai` 实现,通过 `BaseURL` 对接 DashScope:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
chatModel, _ := openai.NewChatModel(ctx, &openai.ChatModelConfig{
|
|||
|
|
APIKey: cfg.AI.LLM.APIKey,
|
|||
|
|
Model: cfg.AI.LLM.Model,
|
|||
|
|
BaseURL: cfg.AI.LLM.Endpoint, // "https://dashscope.aliyuncs.com/compatible-mode/v1"
|
|||
|
|
})
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 2.4 StreamReader
|
|||
|
|
|
|||
|
|
`schema.StreamReader[T]` 是 Eino 的流式数据抽象:
|
|||
|
|
- `sr.Recv()` 读取一帧,`io.EOF` 表示流结束
|
|||
|
|
- `schema.Pipe[T](bufSize)` 创建 `StreamReader` + `StreamWriter` 对
|
|||
|
|
- 框架自动处理 `T ↔ StreamReader[T]` 的转换(装箱/concat)
|
|||
|
|
|
|||
|
|
### 2.5 Callback
|
|||
|
|
|
|||
|
|
Callback 是 Eino 的 AOP 机制,支持节点生命周期钩子:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
type Handler interface {
|
|||
|
|
OnStart(ctx, *RunInfo, CallbackInput) context.Context
|
|||
|
|
OnEnd(ctx, *RunInfo, CallbackOutput) context.Context
|
|||
|
|
OnError(ctx, *RunInfo, error) context.Context
|
|||
|
|
OnStartWithStreamInput(ctx, *RunInfo, *StreamReader[CallbackInput]) context.Context
|
|||
|
|
OnEndWithStreamOutput(ctx, *RunInfo, *StreamReader[CallbackOutput]) context.Context
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
CamTalk 使用 `utils/callbacks.NewHandlerHelper()` 构建 typed handler:
|
|||
|
|
- `ModelCallbackHandler.OnEndWithStreamOutput`:逐 token 推送 `llm_chunk`
|
|||
|
|
|
|||
|
|
### 2.6 State
|
|||
|
|
|
|||
|
|
Graph 全局状态,通过 `WithGenLocalState` 注册:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
type PipelineState struct {
|
|||
|
|
FullResponse strings.Builder
|
|||
|
|
TranscribedText string
|
|||
|
|
TokenUsage *TokenUsage
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
g := compose.NewGraph[I, O](compose.WithGenLocalState(func(ctx context.Context) *PipelineState {
|
|||
|
|
return &PipelineState{}
|
|||
|
|
}))
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
节点通过 `compose.ProcessState` 读写 State。
|
|||
|
|
|
|||
|
|
## 3. CamTalk Graph 设计
|
|||
|
|
|
|||
|
|
### 3.1 拓扑结构
|
|||
|
|
|
|||
|
|
```
|
|||
|
|
START → STT → History → ChatModel → Splitter → TTS → Done → END
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
| 节点 | 类型 | 输入 → 输出 | 职责 |
|
|||
|
|
|------|------|------------|------|
|
|||
|
|
| STT | InvokableLambda | `PipelineInput → STTOutput` | 语音识别,写入 State |
|
|||
|
|
| History | InvokableLambda | `STTOutput → []*schema.Message` | 组装提示词和历史 |
|
|||
|
|
| ChatModel | ChatModel(原生) | `[]*schema.Message → StreamReader[*Message]` | LLM 流式推理 |
|
|||
|
|
| Splitter | TransformableLambda | `StreamReader[string] → StreamReader[[]string]` | 句子切分 |
|
|||
|
|
| TTS | InvokableLambda | `[]string → struct{}` | 语音合成,推送音频 |
|
|||
|
|
| Done | InvokableLambda | `struct{} → PipelineOutput` | 发送 llm_done |
|
|||
|
|
|
|||
|
|
### 3.2 数据类型定义
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// Graph 统一输入
|
|||
|
|
type PipelineInput struct {
|
|||
|
|
AudioData []byte // base64 解码后的音频(可选)
|
|||
|
|
ImageData []byte // base64 解码后的图像(可选)
|
|||
|
|
Text string // 直接文本输入(可选,跳过 STT)
|
|||
|
|
SessionID string
|
|||
|
|
RequestID string
|
|||
|
|
Language string // zh / en
|
|||
|
|
Scenario string // free_chat, interviewer, etc.
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// Graph 统一输出
|
|||
|
|
type PipelineOutput struct {
|
|||
|
|
TranscribedText string // STT 结果
|
|||
|
|
FullResponse string // LLM 完整回复
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// Pipeline State(跨节点共享)
|
|||
|
|
type PipelineState struct {
|
|||
|
|
FullResponse strings.Builder
|
|||
|
|
TranscribedText string
|
|||
|
|
TokenUsage *TokenUsage
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 3.3 流式模式
|
|||
|
|
|
|||
|
|
Graph 使用 **Stream 模式**调用:
|
|||
|
|
- 内部所有节点以 Transform 模式运行
|
|||
|
|
- ChatModel 的 `Stream()` 方法实现真正的 token 级流式
|
|||
|
|
- 适配器消费 `StreamReader[PipelineOutput]` 触发整条链路
|
|||
|
|
|
|||
|
|
### 3.4 消息推送机制
|
|||
|
|
|
|||
|
|
| 消息 | 推送方式 | 时机 |
|
|||
|
|
|------|---------|------|
|
|||
|
|
| `stt_result` | Lambda 内部直接调用 Sender | STT 完成后 |
|
|||
|
|
| `llm_chunk` | Callback `OnEndWithStreamOutput` | ChatModel 逐 token |
|
|||
|
|
| `tts_audio` | Lambda 内部直接调用 Sender | TTS 逐句合成 |
|
|||
|
|
| `llm_done` | Lambda 内部直接调用 Sender | Done 节点执行时 |
|
|||
|
|
|
|||
|
|
**Context 注入**:Sender、RequestID、SessionID、PipelineState 通过 `context.WithValue` 传递。
|
|||
|
|
|
|||
|
|
### 3.5 多模态支持
|
|||
|
|
|
|||
|
|
History 节点将图片构建为 `schema.Message.UserInputMultiContent`:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
systemMsg.UserInputMultiContent = []schema.MessageInputPart{
|
|||
|
|
{
|
|||
|
|
Type: schema.ChatMessagePartTypeImageURL,
|
|||
|
|
Image: &schema.MessageInputImage{
|
|||
|
|
MessagePartCommon: schema.MessagePartCommon{
|
|||
|
|
Base64Data: &base64Str,
|
|||
|
|
MIMEType: "image/jpeg",
|
|||
|
|
},
|
|||
|
|
Detail: schema.ImageURLDetailAuto,
|
|||
|
|
},
|
|||
|
|
},
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
## 4. 实现要点
|
|||
|
|
|
|||
|
|
### 4.1 目录结构
|
|||
|
|
|
|||
|
|
```
|
|||
|
|
backend/internal/eino/
|
|||
|
|
├── types.go # PipelineInput/Output、STTOutput、TokenUsage
|
|||
|
|
├── state.go # PipelineState(跨节点状态)
|
|||
|
|
├── callback.go # Callback handler(LLM token 推送)
|
|||
|
|
├── graph.go # Graph 构建与编译
|
|||
|
|
├── adapter.go # EinoOrchestrator(Orchestrator 接口适配器)
|
|||
|
|
├── nodes_stt.go # STT Lambda
|
|||
|
|
├── nodes_history.go # 历史组装 Lambda
|
|||
|
|
├── nodes_splitter.go # 句子分割 Transform Lambda
|
|||
|
|
├── nodes_tts.go # TTS Lambda
|
|||
|
|
├── nodes_done.go # Done Lambda
|
|||
|
|
└── graph_test.go # 单元测试
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 4.2 关键节点实现
|
|||
|
|
|
|||
|
|
#### STT Lambda(可选跳过)
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func sttLambda(sttSvc stt.Service) func(ctx context.Context, input PipelineInput) (STTOutput, error) {
|
|||
|
|
return func(ctx context.Context, input PipelineInput) (STTOutput, error) {
|
|||
|
|
// 文本模式:跳过 STT
|
|||
|
|
if input.Text != "" {
|
|||
|
|
return STTOutput{Text: input.Text, Language: input.Language}, nil
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 调用 STT 服务
|
|||
|
|
result, err := sttSvc.Recognize(ctx, input.AudioData, stt.Options{
|
|||
|
|
Language: input.Language,
|
|||
|
|
})
|
|||
|
|
if err != nil {
|
|||
|
|
return STTOutput{}, fmt.Errorf("STT error: %w", err)
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
return STTOutput{Text: result.Text, Language: result.Language}, nil
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
#### Splitter Transform Lambda(句子切分)
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func splitterLambda() func(ctx, *schema.StreamReader[*schema.Message]) (*schema.StreamReader[[]string], error) {
|
|||
|
|
return func(ctx context.Context, stream *schema.StreamReader[*schema.Message]) (*schema.StreamReader[[]string], error) {
|
|||
|
|
sr, sw := schema.Pipe[[]string](8)
|
|||
|
|
|
|||
|
|
go func() {
|
|||
|
|
defer sw.Close()
|
|||
|
|
var buffer []rune
|
|||
|
|
|
|||
|
|
for {
|
|||
|
|
chunk, err := stream.Recv()
|
|||
|
|
if err != nil {
|
|||
|
|
if err == io.EOF {
|
|||
|
|
if len(buffer) > 0 {
|
|||
|
|
sw.Send([]string{string(buffer)}, nil)
|
|||
|
|
}
|
|||
|
|
return
|
|||
|
|
}
|
|||
|
|
sw.Send(nil, err)
|
|||
|
|
return
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
for _, r := range chunk.Content {
|
|||
|
|
buffer = append(buffer, r)
|
|||
|
|
if isSentenceDelimiter(r) {
|
|||
|
|
sw.Send([]string{string(buffer)}, nil)
|
|||
|
|
buffer = buffer[:0]
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
}()
|
|||
|
|
|
|||
|
|
return sr, nil
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
#### TTS Lambda(并行合成)
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func ttsLambda(ttsSvc tts.Service, sender orchestrator.Sender) func(ctx, []string) (struct{}, error) {
|
|||
|
|
return func(ctx context.Context, sentences []string) (struct{}, error) {
|
|||
|
|
for _, sentence := range sentences {
|
|||
|
|
if sentence == "" {
|
|||
|
|
continue
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 调用 TTS 服务
|
|||
|
|
audioData, err := ttsSvc.Synthesize(ctx, sentence, tts.Options{})
|
|||
|
|
if err != nil {
|
|||
|
|
// TTS 失败不中断流程,仅记录日志
|
|||
|
|
log.Warn("TTS synthesis failed", zap.Error(err))
|
|||
|
|
continue
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// 推送音频到客户端
|
|||
|
|
sender.SendTTSAudio(orchestrator.TTSAudioPayload{
|
|||
|
|
Audio: audioData,
|
|||
|
|
Format: "mp3",
|
|||
|
|
})
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
return struct{}{}, nil
|
|||
|
|
}
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 4.3 Callback 集成
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// ModelCallbackHandler 用于 LLM token 推送
|
|||
|
|
type ModelCallbackHandler struct {
|
|||
|
|
sender orchestrator.Sender
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
func (h *ModelCallbackHandler) OnEndWithStreamOutput(
|
|||
|
|
ctx context.Context,
|
|||
|
|
info *callbacks.RunInfo,
|
|||
|
|
output *schema.StreamReader[*schema.Message],
|
|||
|
|
) context.Context {
|
|||
|
|
// 逐 token 推送到客户端
|
|||
|
|
for {
|
|||
|
|
msg, err := output.Recv()
|
|||
|
|
if err == io.EOF {
|
|||
|
|
break
|
|||
|
|
}
|
|||
|
|
if err != nil {
|
|||
|
|
return ctx
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
h.sender.SendLLMChunk(orchestrator.LLMChunkPayload{
|
|||
|
|
Content: msg.Content,
|
|||
|
|
})
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
return ctx
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 4.4 按请求动态配置
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// 运行时 Option:每请求可变
|
|||
|
|
func WithModelName(name string) compose.Option {
|
|||
|
|
return compose.WithChatModelOption(model.WithModel(name))
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
func WithTemperature(temp float32) compose.Option {
|
|||
|
|
return compose.WithChatModelOption(model.WithTemperature(temp))
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
// WebSocket Handler 中的调用
|
|||
|
|
func (c *Client) handleQuery(req QueryRequest) {
|
|||
|
|
opts := []compose.Option{}
|
|||
|
|
|
|||
|
|
if req.Model != "" {
|
|||
|
|
opts = append(opts, WithModelName(req.Model))
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
output, err := c.pipeline.Stream(ctx, PipelineInput{...}, opts...)
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 4.5 注意事项
|
|||
|
|
|
|||
|
|
#### 值类型 vs 指针类型
|
|||
|
|
Graph 泛型参数必须使用值类型(`PipelineInput`/`PipelineOutput`),所有 Lambda 的输入输出也使用值类型。框架在 Transform 模式下会自动处理 `T` 和 `StreamReader[T]` 的转换。
|
|||
|
|
|
|||
|
|
#### Callback 运行时传入
|
|||
|
|
Callback 通过 `Stream()` 的 option 传入,不在 `Compile()` 时注册:
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
streamReader, err := runnable.Stream(ctx, input, compose.WithCallbacks(handler))
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
#### eino-ext 与 DashScope 兼容性
|
|||
|
|
eino-ext OpenAI ChatModel 通过 `BaseURL` 对接 DashScope 兼容接口。需注意:
|
|||
|
|
- 多模态图片使用 `Base64Data` + `MIMEType` 格式
|
|||
|
|
- `Timeout` 控制单次请求超时
|
|||
|
|
- 流式输出通过 `Stream()` 方法获取 `StreamReader[*schema.Message]`
|
|||
|
|
|
|||
|
|
#### 框架自动类型转换
|
|||
|
|
Eino 框架在编排场景中自动处理以下转换:
|
|||
|
|
- **T → StreamReader[T]**:将完整值装箱为单帧流(非流式 → 假流式)
|
|||
|
|
- **StreamReader[T] → T**:将流 concat 为完整值(流式 → 非流式)
|
|||
|
|
|
|||
|
|
这使得不同流式模式的节点可以无缝连接。
|
|||
|
|
|
|||
|
|
## 5. 测试策略
|
|||
|
|
|
|||
|
|
### 5.1 单元测试
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
func TestPipelineGraph_WithTextInput(t *testing.T) {
|
|||
|
|
mockLLM := &mockChatModel{responses: []string{"你好!"}}
|
|||
|
|
mockSender := &mockSender{}
|
|||
|
|
|
|||
|
|
graph, err := NewPipelineGraph(ctx, &GraphOption{
|
|||
|
|
ChatModel: mockLLM,
|
|||
|
|
Sender: mockSender,
|
|||
|
|
})
|
|||
|
|
require.NoError(t, err)
|
|||
|
|
|
|||
|
|
output, err := graph.Invoke(ctx, PipelineInput{
|
|||
|
|
Text: "你好",
|
|||
|
|
SessionID: "test-session",
|
|||
|
|
})
|
|||
|
|
require.NoError(t, err)
|
|||
|
|
assert.Equal(t, "你好!", output.FullResponse)
|
|||
|
|
assert.True(t, mockSender.LLMDoneSent)
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
func TestPipelineGraph_WithAudioInput(t *testing.T) {
|
|||
|
|
mockSTT := &mockSTT{text: "你好"}
|
|||
|
|
mockLLM := &mockChatModel{responses: []string{"你好!"}}
|
|||
|
|
mockTTS := &mockTTS{audio: []byte("fake-audio")}
|
|||
|
|
mockSender := &mockSender{}
|
|||
|
|
|
|||
|
|
graph, _ := NewPipelineGraph(ctx, &GraphOption{
|
|||
|
|
ChatModel: mockLLM,
|
|||
|
|
STTService: mockSTT,
|
|||
|
|
TTSService: mockTTS,
|
|||
|
|
Sender: mockSender,
|
|||
|
|
})
|
|||
|
|
|
|||
|
|
output, err := graph.Invoke(ctx, PipelineInput{
|
|||
|
|
AudioData: []byte("fake-audio-data"),
|
|||
|
|
SessionID: "test-session",
|
|||
|
|
})
|
|||
|
|
require.NoError(t, err)
|
|||
|
|
assert.True(t, mockSender.TTSAudioSent)
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 5.2 集成测试
|
|||
|
|
|
|||
|
|
- 启动真实 OpenAI API 调用(使用测试 key)
|
|||
|
|
- 验证 WebSocket 消息序列:`stt_result` → `llm_chunk` × N → `llm_done` → `tts_audio` × N
|
|||
|
|
- 验证 interrupt 取消功能
|
|||
|
|
- 验证多并发请求隔离
|
|||
|
|
|
|||
|
|
## 6. 未来扩展路径
|
|||
|
|
|
|||
|
|
基于 Eino Graph 的重构完成后,可无缝扩展:
|
|||
|
|
|
|||
|
|
1. **ReAct Agent**:Graph 添加 Branch 节点,实现 LLM → Tool → LLM 循环
|
|||
|
|
2. **多模态理解**:添加视觉分析 Lambda 节点(图像描述 → 上下文注入)
|
|||
|
|
3. **Model Router**:Graph 前置分支节点,按场景/成本路由不同 LLM
|
|||
|
|
4. **Rate Limiter**:通过 Callback 的 OnStart 实现令牌桶
|
|||
|
|
5. **Checkpoint/Resume**:利用 Eino 的 CheckpointStore 实现断点续传
|
|||
|
|
6. **Multi-Agent**:利用 ADK 的 Supervisor/SequentialAgent 编排复杂对话流程
|
|||
|
|
|
|||
|
|
## 附录:关键 Eino API 参考
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// 构建 Graph
|
|||
|
|
g := compose.NewGraph[I, O](opts...)
|
|||
|
|
g.AddChatModelNode(key, chatModel)
|
|||
|
|
g.AddLambdaNode(key, lambda, opts...)
|
|||
|
|
g.AddEdge(from, to)
|
|||
|
|
g.AddBranch(from, branchFunc, mapping)
|
|||
|
|
|
|||
|
|
// 编译
|
|||
|
|
runnable, err := g.Compile(ctx, opts...)
|
|||
|
|
|
|||
|
|
// 执行四种模式
|
|||
|
|
output, err := runnable.Invoke(ctx, input, opts...)
|
|||
|
|
stream, err := runnable.Stream(ctx, input, opts...)
|
|||
|
|
output, err := runnable.Collect(ctx, inputStream, opts...)
|
|||
|
|
stream, err := runnable.Transform(ctx, inputStream, opts...)
|
|||
|
|
|
|||
|
|
// Lambda 四种构造器
|
|||
|
|
lambda := compose.InvokableLambda(fn) // I → O
|
|||
|
|
lambda := compose.StreamableLambda(fn) // I → StreamReader[O]
|
|||
|
|
lambda := compose.CollectableLambda(fn) // StreamReader[I] → O
|
|||
|
|
lambda := compose.TransformableLambda(fn) // StreamReader[I] → StreamReader[O]
|
|||
|
|
|
|||
|
|
// Stream 操作
|
|||
|
|
sr, sw := schema.Pipe[T](bufSize)
|
|||
|
|
sw.Send(chunk, err)
|
|||
|
|
chunk, err := sr.Recv()
|
|||
|
|
sw.Close()
|
|||
|
|
|
|||
|
|
// Option
|
|||
|
|
compose.WithCallbacks(handler)
|
|||
|
|
compose.WithCallbacks(handler).DesignateNode("node_key")
|
|||
|
|
compose.WithChatModelOption(model.WithTemperature(0.7))
|
|||
|
|
compose.WithGenLocalState(genFunc)
|
|||
|
|
```
|