## 主要变更 ### 文档重构(减少 1199 行,-23%) - 01-架构设计.md: 503→369 行 (-27%),删除 DDL/配置示例,精简鉴权/存储描述 - 02-接口文档.md: 1313→570 行 (-57%),删除 Go 接口/Orchestrator 实现/配置管理 - 07-成本控制.md: 65→59 行 (-9%),代码块替换为文件引用 ### 文档编号规范化 - 08-功能创意.md → 删除(内容整合到 README.md "功能扩展方向") - 10-Eino框架与编排设计.md → 08-Eino框架与编排设计.md - 情景切换.md → 09-情景切换.md - 12-鉴权体系.md → 10-鉴权体系.md - 13-令牌桶限流.md → 11-令牌桶限流.md ### 交叉引用更新 - 01-架构设计.md: 更新对鉴权体系/令牌桶限流的引用为新编号 - README.md: 更新文档索引表、推荐阅读顺序、新增功能扩展方向 ### 删除过时文档 - 09-技术名词解释.md(内容已整合到 03-技术选型.md) - 10-Eino重构方案.md(历史记录,已完成) - 11-Eino框架技术文档.md(已合并到 08) - 情景切换功能完整文档.md(已规范化为 09) ## 重构原则 - 架构文档聚焦系统结构,移除实现细节 - 接口文档保留纯契约,删除内部实现 - 编号连续(01-11),语义清晰 - 通过交叉引用连接相关文档,避免重复
117 lines
3.4 KiB
Go
117 lines
3.4 KiB
Go
package eino
|
||
|
||
import (
|
||
"context"
|
||
"time"
|
||
|
||
openaiImpl "github.com/cloudwego/eino-ext/components/model/openai"
|
||
"github.com/cloudwego/eino/compose"
|
||
|
||
"github.com/hhs/camtalk/internal/ai/stt"
|
||
"github.com/hhs/camtalk/internal/ai/tts"
|
||
"github.com/hhs/camtalk/internal/config"
|
||
"github.com/hhs/camtalk/internal/logger"
|
||
"github.com/hhs/camtalk/internal/models"
|
||
"github.com/hhs/camtalk/internal/session"
|
||
)
|
||
|
||
const (
|
||
nodeSTT = "stt"
|
||
nodeHistory = "history"
|
||
nodeLLM = "llm"
|
||
nodeMessageToString = "msg2str"
|
||
nodeSplitter = "splitter"
|
||
nodeTTS = "tts"
|
||
nodeDone = "done"
|
||
)
|
||
|
||
// PipelineGraph 封装编译后的 Eino Graph。
|
||
type PipelineGraph struct {
|
||
Runnable compose.Runnable[PipelineInput, PipelineOutput]
|
||
}
|
||
|
||
// NewPipelineGraph 构建 CamTalk AI 编排 Graph。
|
||
//
|
||
// 拓扑:START → STT → History → ChatModel → Splitter → TTS → Done → END
|
||
//
|
||
// Graph 使用 Stream 模式调用,ChatModel 实现真正的 token 级流式输出。
|
||
// LLM token 通过 Callback 的 OnEndWithStreamOutput 实时推送到客户端。
|
||
func NewPipelineGraph(
|
||
ctx context.Context,
|
||
cfg *config.Config,
|
||
sttService stt.Service,
|
||
ttsService tts.Service,
|
||
sessionMgr session.Manager,
|
||
) (*PipelineGraph, error) {
|
||
log := logger.Log
|
||
|
||
// 1. 创建 eino-ext ChatModel(对接 DashScope OpenAI 兼容接口)
|
||
chatModel, err := openaiImpl.NewChatModel(ctx, &openaiImpl.ChatModelConfig{
|
||
APIKey: cfg.AI.LLM.APIKey,
|
||
Model: cfg.AI.LLM.Model,
|
||
BaseURL: cfg.AI.LLM.Endpoint,
|
||
Timeout: time.Duration(cfg.AI.LLM.Timeout) * time.Second,
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
log.Infow("Eino ChatModel 初始化成功",
|
||
"model", cfg.AI.LLM.Model,
|
||
"endpoint", cfg.AI.LLM.Endpoint)
|
||
|
||
// 2. 构建 Graph(值类型,非指针)
|
||
g := compose.NewGraph[PipelineInput, PipelineOutput](
|
||
compose.WithGenLocalState(genLocalState),
|
||
)
|
||
|
||
// 3. 添加节点
|
||
maxHistory := cfg.Session.MaxHistory
|
||
|
||
_ = g.AddLambdaNode(nodeSTT, NewSTTLambda(sttService))
|
||
_ = g.AddLambdaNode(nodeHistory, NewHistoryLambda(sessionMgr.GetHistory, maxHistory))
|
||
_ = g.AddChatModelNode(nodeLLM, chatModel)
|
||
_ = g.AddLambdaNode(nodeMessageToString, NewMessageToStringLambda())
|
||
_ = g.AddLambdaNode(nodeSplitter, NewSplitterLambda())
|
||
_ = g.AddLambdaNode(nodeTTS, NewTTSLambda(
|
||
ttsService,
|
||
cfg.AI.TTS.Voice,
|
||
cfg.AI.TTS.Speed,
|
||
cfg.AI.TTS.OutputFormat,
|
||
cfg.AI.TTS.SampleRate,
|
||
))
|
||
_ = g.AddLambdaNode(nodeDone, NewDoneLambda(cfg.AI.LLM.Model))
|
||
|
||
// 4. 连接边
|
||
_ = g.AddEdge(compose.START, nodeSTT)
|
||
_ = g.AddEdge(nodeSTT, nodeHistory)
|
||
_ = g.AddEdge(nodeHistory, nodeLLM)
|
||
_ = g.AddEdge(nodeLLM, nodeMessageToString)
|
||
_ = g.AddEdge(nodeMessageToString, nodeSplitter)
|
||
_ = g.AddEdge(nodeSplitter, nodeTTS)
|
||
_ = g.AddEdge(nodeTTS, nodeDone)
|
||
_ = g.AddEdge(nodeDone, compose.END)
|
||
|
||
// 5. 编译(回调在运行时通过 Stream option 传入)
|
||
runnable, err := g.Compile(ctx)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
log.Infow("Eino Graph 编译成功", "nodes", 7)
|
||
return &PipelineGraph{Runnable: runnable}, nil
|
||
}
|
||
|
||
// buildPipelineInput 从 WebSocket 请求和会话配置构建 Graph 输入。
|
||
func buildPipelineInput(req models.WsQuery, sessionID string, sess *models.Session, audioData, imageData []byte) PipelineInput {
|
||
return PipelineInput{
|
||
AudioData: audioData,
|
||
ImageData: imageData,
|
||
Text: req.Text,
|
||
SessionID: sessionID,
|
||
RequestID: req.RequestID,
|
||
Language: sess.Config.Language,
|
||
Scenario: sess.Config.Scenario,
|
||
TTSEnabled: sess.Config.TTSEnabled,
|
||
}
|
||
}
|