diff --git a/CLAUDE.md b/CLAUDE.md index 8a8a26a..9798898 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -4,7 +4,7 @@ ## 项目概述 -CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头和麦克风与 AI 交互,AI 理解视觉场景和语音输入后,以文字和语音形式给出自然回应。项目目前处于设计文档阶段,源代码正在逐步构建。 +CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头和麦克风与 AI 交互,AI 理解视觉场景和语音输入后,以文字和语音形式给出自然回应。 > **文档优先原则:** 执行任何开发任务前,先读取 `docs/` 下的相关设计文档(架构、接口、技术选型等),以文档为最高依据。代码实现应与文档一致;若有偏差,优先更新文档(尤其是接口文档)。 @@ -13,12 +13,12 @@ CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头 三层系统: 1. **浏览器客户端**(React 18 + TypeScript, Vite)—— 媒体采集、边缘预处理(VAD 通过 `@ricky0123/vad-web`、关键帧检测通过 Canvas 像素比较)、UI 渲染。核心 Hook:`useVisionSession()` -2. **Go 网关**(Gin, gorilla/websocket, Viper, Zap)—— WebSocket 服务器、会话管理、AI 编排。每个 WebSocket 连接一个 goroutine。 -3. **云端 AI 服务** —— 通过 OpenAI 兼容接口可灵活切换。默认:GPT-4o(LLM)、Deepgram(STT)、OpenAI TTS。仅通过 Go 网关访问,浏览器不直连。 +2. **Go 网关**(Gin, gorilla/websocket, Viper, Zap)—— WebSocket 服务器、会话管理、AI 编排(基于 CloudWeGo Eino Graph)。每个 WebSocket 连接一个 goroutine。 +3. **云端 AI 服务** —— 通过 OpenAI 兼容接口可灵活切换。默认:DashScope qwen3-vl-plus(LLM)、MiMo ASR(STT)、MiMo TTS(TTS)。仅通过 Go 网关访问,浏览器不直连。 -**关键模式**:LLM 文本流和 TTS 音频流并行推送给客户端,以最小化感知延迟。 +**关键模式**:AI 编排基于 Eino Graph 声明式 DAG(`START → STT → History → ChatModel → Msg2Str → Splitter → TTS → Done → END`),LLM token 通过 Callback 实时推送,TTS 逐句合成并行推送,最小化感知延迟。 -**存储**:MVP 阶段使用进程内存(`MemoryManager`),Redis 实现已就绪可通过配置切换,PostgreSQL 为规划中。Repository 接口模式(`HistoryRepository`、`UsageRepository`),MVP 用内存实现。 +**存储**:三级存储架构(TieredManager)—— L1 Memory → L2 Redis → L3 PostgreSQL,自动降级。Repository 接口模式(UserRepository、MessageRepository、SessionRepository),PostgreSQL + 内存双实现。 ## 技术栈 @@ -26,9 +26,10 @@ CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头 |------|------| | 前端 | React 18, TypeScript, Vite, @ricky0123/vad-web | | 后端 | Go, Gin, gorilla/websocket, Viper, Zap | -| LLM | GPT-4o(默认,通过 OpenAI 兼容接口可切换) | -| STT | Deepgram(默认) / MiMo ASR | -| TTS | OpenAI TTS(默认) / MiMo TTS | +| AI 编排 | CloudWeGo Eino Graph(声明式 DAG 编排) | +| LLM | DashScope qwen3-vl-plus(默认,通过 eino-ext OpenAI ChatModel 接入) | +| STT | MiMo ASR(默认) / Deepgram | +| TTS | MiMo TTS(默认) / OpenAI TTS | ## 构建与运行命令 @@ -49,13 +50,13 @@ go test -run TestName ./path # 运行单个测试 go vet ./... # 静态分析 ``` -基础设施:MVP 使用进程内存管理会话状态。Redis 已实现可通过配置切换,PostgreSQL 为规划中。 +基础设施:三级存储架构(L1 Memory → L2 Redis → L3 PostgreSQL),通过配置控制启用层级。 ## WebSocket 协议 -端点:`ws://localhost:8080/ws` +端点:`ws://localhost:8080/ws?token=&conversation_id=` -所有消息为 JSON 文本帧,统一信封格式 `{type, request_id?, timestamp?}`。完整契约见 `docs/03-接口文档.md`。 +所有消息为 JSON 文本帧,统一信封格式 `{type, request_id?, timestamp?}`。完整契约见 `docs/02-接口文档.md`。 **客户端 → 服务端**:`query`(图像 Base64 + 音频 Base64)、`config`、`interrupt`、`ping` **服务端 → 客户端**:`connected`、`stt_result`、`llm_chunk`、`llm_done`、`tts_audio`、`error`、`pong` @@ -66,34 +67,50 @@ go vet ./... # 静态分析 ## REST API(辅助) - `GET /api/health` — 健康检查(版本、运行时间、活跃会话数) -- `POST /api/sessions` — 创建会话(可选,MVP 在 WS 连接时自动创建) -- `DELETE /api/sessions/{id}` — 销毁会话 +- `POST /api/auth/register` — 注册 +- `POST /api/auth/login` — 登录 +- `POST /api/auth/refresh` — 刷新 Token +- `POST /api/auth/logout` — 登出 +- `GET /api/conversations` — 对话列表 +- `POST /api/conversations` — 创建对话 +- `GET/PUT/PATCH/DELETE /api/conversations/:id` — 对话 CRUD +- `GET /api/conversations/:id/messages` — 获取对话消息 ## 错误码 -`INVALID_MESSAGE`、`SESSION_NOT_FOUND`、`RATE_LIMITED`、`IMAGE_TOO_LARGE`、`AUDIO_TOO_SHORT`、`LLM_TIMEOUT`、`LLM_ERROR`、`STT_ERROR`、`TTS_ERROR`、`INTERNAL_ERROR` +`INVALID_MESSAGE`、`SESSION_NOT_FOUND`、`RATE_LIMITED`、`IMAGE_TOO_LARGE`、`AUDIO_TOO_SHORT`、`LLM_TIMEOUT`、`LLM_ERROR`、`STT_ERROR`、`TTS_ERROR`、`INTERNAL_ERROR`、`USERNAME_TAKEN`、`INVALID_CREDENTIALS`、`INVALID_TOKEN`、`INVALID_INPUT` ## 前端组件结构 | 组件 | 职责 | |------|------| +| `AuthPage` | 登录/注册表单 | | `CameraManager` | 摄像头流采集 | | `MicManager` | 麦克风音频采集 | | `EdgeProcessor` | VAD + 关键帧检测(Canvas 像素比较) | | `WebSocketManager` | WebSocket 连接生命周期管理 | -| `ChatPanel` | 消息展示 | +| `ChatPanel` | 消息展示、流式回复、文本输入、场景选择 | | `VideoPreview` | 摄像头画面预览 | +| `SessionSidebar` | 左侧抽屉式对话列表 | +| `ConfigPanel` | 右侧抽屉式配置面板 | +| `Toast` | 轻量通知提示 | + +核心 Hook:`useVisionSession()` 封装一次完整的视觉对话会话。 ## 后端模块结构 | 模块 | 职责 | |------|------| -| WebSocket Handler | 连接管理、单播消息推送 | -| Session Manager | 会话状态、对话历史(Memory/Redis,30 分钟 TTL) | -| AI Orchestrator | STT→LLM→TTS 流式并行管道编排 | -| AI Service Layer | AI 服务抽象层(STT/LLM/TTS 多 provider) | -| REST API | 健康检查、会话管理(Gin 路由) | +| WebSocket Handler | 连接管理、JWT 认证、单播消息推送 | +| Session Manager | 会话状态、对话历史(三级存储:Memory/Redis/PostgreSQL,30 分钟 TTL) | +| Eino 编排层 | 基于 Eino Graph 的声明式 AI 编排(7 节点 DAG,Stream 模式,Callback AOP) | +| AI Orchestrator | `EinoOrchestrator` 适配器,包装 Graph 实现 `Orchestrator` 接口 | +| AI Service Layer | AI 服务抽象层(STT/TTS 多 provider,LLM 通过 eino-ext ChatModel) | +| Auth | JWT 双 token 轮转认证,bcrypt 密码哈希 | +| Store | 持久化存储层(UserRepository/MessageRepository/SessionRepository,内存 + PostgreSQL) | +| REST API | 健康检查、认证、对话管理(Gin 路由) | | Models | 数据模型定义 | +| Migrations | 数据库版本化迁移(嵌入式 SQL) | | Model Router | 按请求选择 AI 模型(规划中) | | Rate Limiter | 按用户的令牌桶速率限制(规划中) | diff --git a/docs/01-架构设计.md b/docs/01-架构设计.md index f770b43..02826e4 100644 --- a/docs/01-架构设计.md +++ b/docs/01-架构设计.md @@ -27,7 +27,7 @@ graph TB subgraph Gateway["Go 网关"] WS["WebSocket Handler
连接管理 / 消息分发"] Session["Session Manager
会话状态 / 对话历史"] - Orch["AI Orchestrator
STT→LLM→TTS 流式并行"] + Orch["AI Orchestrator
Eino Graph 声明式编排"] Auth["Auth 模块
JWT / bcrypt"] REST["REST API
健康检查 / 对话管理"] Store["Store 层
Repository 接口"] @@ -70,34 +70,38 @@ graph TB ```mermaid sequenceDiagram participant B as 浏览器 - participant G as Go 网关 + participant G as Go 网关(Eino Graph) participant S as STT - participant L as LLM + participant L as LLM(ChatModel) participant T as TTS B->>B: VAD 检测到语音结束 B->>G: query {image, audio} - G->>S: 音频流 - S-->>G: 流式文本 + Note over G: EinoOrchestrator 启动 Graph.Stream() + + G->>S: STT Lambda:音频 → 文本 + S-->>G: 识别文本 G-->>B: stt_result {text} - G->>L: [图像 + 文本 + 上下文] - loop LLM 流式输出 + G->>G: History Lambda:组装提示词 + 历史 + 多模态消息 + + G->>L: ChatModel Node:流式推理 + loop LLM 流式输出(Callback OnEndWithStreamOutput) L-->>G: token delta G-->>B: llm_chunk {delta} end - G-->>B: llm_done {full_text, tokens} - par LLM 输出的同时 - G->>G: 句子切分器检测到完整句子 - G->>T: 句子文本 - T-->>G: 音频 chunk - G-->>B: tts_audio {audio} - end + G->>G: Msg2Str + Splitter Lambda:句子切分 + G->>T: TTS Lambda:逐句合成 + T-->>G: 音频 chunk + G-->>B: tts_audio {audio} + + G->>G: Done Lambda:发送完成通知 + G-->>B: llm_done {full_text, tokens} G-->>B: tts_audio {final: true} ``` -**关键优化**:LLM 文本流和 TTS 音频流**并行推送**——客户端先逐 token 展示文字,同时 TTS 逐句子合成并推送音频,用户感知延迟大幅降低。 +**关键优化**:Eino Graph 以 Stream 模式运行,ChatModel 的 token 流通过 Callback 的 `OnEndWithStreamOutput` 实时推送到客户端(`llm_chunk`),同时 Splitter 节点将 token 流切分为句子,TTS 节点逐句合成并推送音频。LLM 文本流和 TTS 音频流**并行推送**,用户感知延迟大幅降低。 ## 技术栈 @@ -118,7 +122,8 @@ sequenceDiagram | 语言 | Go | 高并发 goroutine 模型,适合长连接管理 | | HTTP 框架 | Gin | 高性能 HTTP 路由,中间件生态成熟 | | WebSocket | gorilla/websocket | Go 生态最成熟的 WebSocket 库 | -| 会话存储 | Memory(默认) / Redis | 进程内存零依赖,Redis 支持多实例部署 | +| 会话存储 | Memory / Redis / PostgreSQL 三级存储 | 进程内存零依赖,Redis 支持多实例,PG 持久化。TieredManager 自动降级 | +| AI 编排 | CloudWeGo Eino Graph | 声明式 DAG 编排,Stream 模式,Callback AOP | | 持久化存储 | PostgreSQL | 对话历史、用户数据、会话元数据 | | 配置管理 | Viper + godotenv | 支持 YAML + .env + 环境变量覆盖 | | 日志 | Zap | 高性能结构化日志 | @@ -127,11 +132,11 @@ sequenceDiagram | 能力 | 默认方案 | 备选方案 | |------|---------|---------| -| 多模态 LLM | GPT-4o | 通义千问等 OpenAI 兼容模型 | -| 语音识别 STT | Deepgram | MiMo ASR(小米) | -| 语音合成 TTS | OpenAI TTS | MiMo TTS(小米) | +| 多模态 LLM | DashScope qwen3-vl-plus | GPT-4o 等 OpenAI 兼容模型 | +| 语音识别 STT | MiMo ASR(小米) | Deepgram | +| 语音合成 TTS | MiMo TTS(小米) | OpenAI TTS | -> Go 网关的 AI 服务层统一封装不同服务商的调用接口,通过配置切换 provider。 +> Go 网关的 AI 服务层统一封装不同服务商的调用接口,通过配置切换 provider。LLM 通过 Eino 框架的 `eino-ext/components/model/openai` 组件接入,支持任何 OpenAI 兼容接口。 ## 后端模块 @@ -148,14 +153,20 @@ graph LR subgraph Business["业务层"] SM["Session Manager
会话生命周期"] - ORCH["Orchestrator
STT→LLM→TTS 编排"] + ORCH["EinoOrchestrator
Eino Graph 编排"] AS["Auth Service
注册/登录/刷新/登出"] end + subgraph Eino_Layer["Eino 编排层"] + PG["PipelineGraph
7 节点 DAG"] + CB["Callback Handler
LLM token 推送"] + ST["PipelineState
跨节点状态"] + end + subgraph AI_Layer["AI 服务层"] - STT_S["STT Service
Deepgram / MiMo"] - LLM_S["LLM Service
OpenAI 兼容"] - TTS_S["TTS Service
OpenAI / MiMo"] + STT_S["STT Service
MiMo / Deepgram"] + LLM_S["ChatModel
eino-ext OpenAI 兼容"] + TTS_S["TTS Service
MiMo / OpenAI"] end subgraph Data["数据层"] @@ -174,9 +185,12 @@ graph LR WSH --> ORCH APH --> SM APH --> AS - ORCH --> STT_S - ORCH --> LLM_S - ORCH --> TTS_S + ORCH --> PG + PG --> CB + PG --> ST + PG --> STT_S + PG --> LLM_S + PG --> TTS_S SM --> MR SM --> SR AS --> UR @@ -186,7 +200,8 @@ graph LR |------|------| | WebSocket Handler | 管理客户端连接生命周期,JWT 认证,conversation_id 恢复,单播消息推送 | | Session Manager | 维护用户会话状态、对话历史。Memory(默认)/ Redis(可切换),30 分钟 TTL,Write-Through 到 PG | -| AI Orchestrator | 编排 STT→LLM→TTS 流式并行管道,context 取消 + 超时控制 + 句子切分 | +| Eino 编排层 | 基于 CloudWeGo Eino Graph 的声明式 AI 编排。7 节点 DAG(STT→History→ChatModel→Msg2Str→Splitter→TTS→Done),Stream 模式调用,Callback 实现 LLM token 实时推送 | +| AI Orchestrator | `EinoOrchestrator` 适配器,包装 Eino Graph 实现 `Orchestrator` 接口。context 取消 + 超时控制 | | AI Service Layer | AI 服务抽象层,多 provider 支持(Deepgram/MiMo/OpenAI 等) | | Auth | 用户认证与授权。JWT (HS256) 双 token 轮转,bcrypt 密码哈希,Gin 中间件 | | Store | 持久化存储层。UserRepository / MessageRepository / SessionRepository,内存 + PostgreSQL 双实现 | @@ -214,6 +229,41 @@ graph LR 核心 Hook:`useVisionSession()` 封装一次完整的视觉对话会话(摄像头、VAD、WebSocket、消息状态、认证、场景模式)。 +### 前端会话状态模型(三态) + +前端 UI 存在三个会话状态,由 `isConnected` 和 `isCameraOn` 联合决定: + +``` +┌──────────┐ startSession() ┌──────────┐ +│ initial │ ──────────────────→ │ video │ +│ 初始态 │ │ 视频通话 │ +└──────────┘ └──────────┘ + ↑ │ + │ stopSession() stopVideo() + │ │ + │ ▼ + │ ┌──────────┐ + └──────────────────────── │ textOnly │ + │ 文字对话 │ + └──────────┘ + │ + startSession() + │ + ▼ + ┌──────────┐ + │ video │ + └──────────┘ +``` + +| 状态 | 条件 | WebSocket | 摄像头 | 消息 | 文字输入 | +|------|------|-----------|--------|------|---------| +| `initial` | `!isConnected && messages.length === 0` | 断开 | 关闭 | 空 | 可用(自动连接) | +| `video` | `isConnected && isCameraOn` | 连接 | 开启 | 有 | 可用 | +| `textOnly` | `isConnected && !isCameraOn` | 连接 | 关闭 | 保留 | 可用 | + +- **`stopVideo()`**:停止摄像头/麦克风/VAD,保持 WebSocket 连接和消息历史,用户可继续文字对话 +- **`stopSession()`**:完全断开 WebSocket、清空消息、重置状态,回到初始态 + ## 数据库设计 ### ER 关系 @@ -306,10 +356,23 @@ CREATE TABLE refresh_tokens ( | 场景 | 存储方案 | 说明 | |------|---------|------| | 默认 | Memory(进程内) | 零依赖,快速启动。MemoryManager 支持 Write-Through 到 PG | -| 持久化 | Memory + PostgreSQL | 通过 `storage.driver: postgres` 启用,MemoryManager 注入 PG Repository | +| 持久化 | Memory + PostgreSQL | 通过 `storage.persistence.enabled: true` 启用,MemoryManager 注入 PG Repository | | 多实例 | Redis(独立) | 通过配置切换到 RedisManager,适合多实例部署 | +| 三级存储 | TieredManager | L1 Memory → L2 Redis → L3 PostgreSQL,自动降级 | -冷热分离:Redis/Memory 存"热数据"(当前对话上下文,微秒级读写),PostgreSQL 存"冷数据"(历史记录)。MemoryManager 的 Write-Through 机制确保每次 AppendMessage 同时写入 PG,重启后可从 PG 恢复会话。 +**三级存储架构**(`TieredManager`): + +``` +TieredManager +├── L1: Memory(进程内缓存,微秒级读写) +├── L2: Redis(分布式缓存,毫秒级读写) +└── L3: PostgreSQL(持久化存储,冷数据) +``` + +- **读取路径**:L1 → L2 → L3,逐级回源,命中后向上回填 +- **写入路径**:L1 → L2(同步) → L3(异步) +- **健康检查**:后台 goroutine 每 30 秒 ping Redis,故障时自动降级为 L1+L3 模式 +- **冷热分离**:L1/L2 存"热数据"(当前对话上下文),L3 存"冷数据"(历史记录) ## 认证设计 diff --git a/docs/02-接口文档.md b/docs/02-接口文档.md index 071a75e..820a0c4 100644 --- a/docs/02-接口文档.md +++ b/docs/02-接口文档.md @@ -643,44 +643,49 @@ type Options struct { | Provider | 连接方式 | 说明 | |----------|---------|------| -| Deepgram(默认) | WebSocket `wss://api.deepgram.com/v1/listen` | 流式识别,延迟极低,模型 nova-2 | -| MiMo ASR | HTTP POST OpenAI 兼容 `/chat/completions` | 国产替代,PCM 自动转 WAV,支持 zh/en/auto | +| MiMo ASR(默认) | HTTP POST OpenAI 兼容 `/chat/completions` | 国产替代,PCM 自动转 WAV,支持 zh/en/auto | +| Deepgram | WebSocket `wss://api.deepgram.com/v1/listen` | 流式识别,延迟极低,模型 nova-2 | ### LLM 服务接口 -多模态推理:接收图像 + 文本 + 对话历史,流式返回回复。 +多模态推理通过 Eino 框架的 `eino-ext/components/model/openai` ChatModel 组件实现,替代了原有的手动 `llm.Service` 接口。 + +**Eino ChatModel 配置**: ```go -// Service 多模态大模型服务契约。 -type Service interface { - // ChatStream 流式推理,返回增量文本的 channel。 - ChatStream(ctx context.Context, req Request) (<-chan Chunk, error) -} +chatModel, _ := openaiImpl.NewChatModel(ctx, &openaiImpl.ChatModelConfig{ + APIKey: cfg.AI.LLM.APIKey, + Model: cfg.AI.LLM.Model, // 默认 "qwen3-vl-plus" + BaseURL: cfg.AI.LLM.Endpoint, // 默认 DashScope OpenAI 兼容接口 + Timeout: time.Duration(cfg.AI.LLM.Timeout) * time.Second, +}) +``` -// Request 推理请求。 -type Request struct { - Image []byte // JPEG 图片(已从 Base64 解码) - Text string // 用户语音识别后的文本 - History []models.Message // 最近 N 轮对话历史 - Language string // "zh-CN" - SystemPrompt string // 系统提示词(含场景 prompt) -} +**ChatModel 接口**(Eino 组件标准接口): -// Chunk 流式推理的一个增量片段。 -type Chunk struct { - Delta string - Done bool - TokensUsed *TokenUsage // 仅 Done=true 时有值 - Model string // 实际使用的模型名 +```go +type BaseChatModel interface { + Generate(ctx, []*schema.Message, ...Option) (*schema.Message, error) + Stream(ctx, []*schema.Message, ...Option) (*schema.StreamReader[*schema.Message], error) } ``` +CamTalk 使用 `Stream()` 模式,通过 Eino Graph 的 Stream 调用触发,token 级流式输出通过 Callback `OnEndWithStreamOutput` 推送到客户端。 + **接入约定**: -- 端点:`POST {endpoint}/chat/completions`,通过配置切换 -- 图片传入:`image_url` 字段使用 `data:image/jpeg;base64,...` 格式 -- 流式响应:`stream: true`,通过 SSE 逐 chunk 返回 -- 超时:10 秒,超时返回 `LLM_TIMEOUT` 错误 -- 系统提示词:根据语言和场景(scenario)动态构建 +- 端点:通过 `ai.llm.endpoint` 配置,支持任何 OpenAI 兼容接口 +- 默认模型:`qwen3-vl-plus`(DashScope),通过 `ai.llm.model` 配置切换 +- 图片传入:History 节点构建 `schema.Message.UserInputMultiContent`,使用 `Base64Data` + `MIMEType` 格式 +- 流式响应:Eino 框架原生 `StreamReader` 支持 +- 超时:通过 `ChatModelConfig.Timeout` 控制 +- 系统提示词:History 节点根据语言和场景(scenario)动态构建 + +**保留的类型定义**(`ai/llm/llm.go`): + +```go +// Request / Chunk / TokenUsage 类型定义仍保留在 ai/llm 包中, +// 供 prompt.go 和 scenarios.go 使用。LLM 推理本身通过 eino-ext ChatModel 执行。 +``` ### TTS 服务接口 @@ -708,16 +713,18 @@ type Options struct { | Provider | 端点 | 说明 | |----------|------|------| -| OpenAI TTS(默认) | `POST /audio/speech` | 逐句合成,返回 MP3 流 | -| MiMo TTS | `POST /chat/completions` | 国产替代,base64 音频响应 | +| MiMo TTS(默认) | `POST /chat/completions` | 国产替代,base64 音频响应 | +| OpenAI TTS | `POST /audio/speech` | 逐句合成,返回 MP3 流 | --- ## 四、AI 编排器(Orchestrator) -### 编排策略:句子级流式并行 +### 编排架构:Eino Graph 声明式编排 -核心矛盾:LLM 流式输出逐 token,TTS 需要完整句子才能合成。解法:**句子切分器 + 管道并行**。 +AI 编排层基于 [CloudWeGo Eino](https://github.com/cloudwego/eino) 框架的 `compose.Graph` 实现,替代了原有的手写 goroutine 管道。Eino Graph 是一个声明式的有向无环图(DAG)编排器,支持类型安全的流式数据传递和 Callback AOP 机制。 + +**核心矛盾**:LLM 流式输出逐 token,TTS 需要完整句子才能合成。解法:**Eino TransformableLambda 句子切分 + Stream 模式管道**。 ``` LLM 流式输出: "这" "是一" "朵红色" "的花。" "它看起" "来很美" "丽。" @@ -732,9 +739,29 @@ LLM 流式输出: "这" "是一" "朵红色" "的花。" "它看起" "来很美 ``` **时序保证**: -- `llm_chunk` 消息一定先于对应句子的 `tts_audio` 到达客户端 +- `llm_chunk` 通过 Callback `OnEndWithStreamOutput` 实时推送,一定先于对应句子的 `tts_audio` 到达客户端 - 用户先看到文字,紧接着听到语音(感知延迟 < 0.5 秒) +### Graph 拓扑 + +``` +START → STT → History → ChatModel → Msg2Str → Splitter → TTS → Done → END +``` + +| 节点 | Lambda 类型 | 输入 → 输出 | 职责 | +|------|------------|------------|------| +| STT | InvokableLambda | `PipelineInput → STTOutput` | 语音识别(文本模式跳过),发送 `stt_result`,写入 State | +| History | InvokableLambda | `STTOutput → []*schema.Message` | 组装系统提示词 + 对话历史 + 多模态图像消息 | +| ChatModel | ChatModel(原生) | `[]*schema.Message → StreamReader[*Message]` | Eino 原生 LLM 流式推理 | +| Msg2Str | TransformableLambda | `StreamReader[*Message] → StreamReader[string]` | 提取 LLM 输出文本 | +| Splitter | TransformableLambda | `StreamReader[string] → StreamReader[string]` | 按句子分隔符切分,逐句输出 | +| TTS | TransformableLambda | `StreamReader[string] → StreamReader[struct{}]` | 逐句调用 TTS 服务,推送 `tts_audio` | +| Done | InvokableLambda | `struct{} → PipelineOutput` | 发送 `llm_done`,返回最终输出 | + +**依赖版本**: +- `github.com/cloudwego/eino v0.9.9` +- `github.com/cloudwego/eino-ext/components/model/openai v0.1.13` + ### Orchestrator 接口 ```go @@ -752,30 +779,76 @@ type Sender interface { } ``` -**Pipeline 实现流程**: -1. Base64 解码音频/图片 -2. 调用 `stt.Recognize()` → 发送 `stt_result` -3. 调用 `llm.ChatStream()` 获取流式输出,goroutine 消费 token → 发送 `llm_chunk` + 句子切分 -4. 另一 goroutine 从句子 channel 读取 → 调用 `tts.SynthesizeStream()` → 发送 `tts_audio` -5. 流结束 → 发送 `llm_done` -6. TTS 失败静默跳过,STT/LLM 失败发送对应 error 消息 +WS Handler 通过 `Orchestrator` 接口与编排层交互,不感知 Eino 实现细节。 + +### EinoOrchestrator 执行流程 + +`EinoOrchestrator` 实现 `Orchestrator` 接口,包装 Eino Graph: + +1. 设置活跃请求,获取会话配置 +2. Base64 解码音频/图片 +3. 构建 `PipelineInput` +4. 注入 context 值(Sender、RequestID、SessionID、PipelineState、StartTime) +5. 追加用户消息到历史 +6. 调用 `graph.Runnable.Stream(ctx, input, callbacks)` — Stream 模式触发整条链路惰性执行 +7. 消费 `StreamReader[PipelineOutput]` 直到 EOF +8. 追加助手消息到历史 + +### Callback 机制 + +LLM token 推送通过 Eino Callback 实现,而非在 Lambda 节点中硬编码: + +```go +// 构建 typed callback handler +handler := callbacks.NewHandlerHelper().ChatModel(&modelCallbackHandler{}).Handler() + +// 运行时传入(不在 Compile 时注册) +streamReader, err := runnable.Stream(ctx, input, compose.WithCallbacks(handler)) +``` + +**OnEndWithStreamOutput** 回调: +- 接收 ChatModel 的 `StreamReader[*schema.Message]` +- 逐 chunk 推送 `llm_chunk` 到客户端 +- 累积完整回复到 `PipelineState` +- 记录 token 用量 + +### State 机制 + +`PipelineState` 是 Graph 级别的线程安全状态,通过 `compose.WithGenLocalState` 注册: + +```go +type PipelineState struct { + FullResponse strings.Builder // LLM 完整回复(Callback 累积) + TranscribedText string // STT 识别文本 + TokenUsage *TokenUsage // Token 用量 + SessionID string + RequestID string + ImageData []byte + Scenario string + Language string + TTSEnabled bool +} +``` + +各节点通过 `stateFromCtx(ctx)` 读写 State,实现跨节点数据共享。 ### 并发控制 - 每个 `ProcessQuery` 调用在独立 goroutine 中运行 -- `context.WithTimeout` 确保 10 秒总超时 -- `interrupt` 消息触发 `cancel()`,LLM/TTS 流式全部中断 +- `context.WithTimeout` 确保总超时 +- `interrupt` 消息触发 `cancel()`,Eino Graph 内部所有流式节点中断 - 同一 session 内同时只允许一个活跃请求,新请求自动取消上一个 +- `PipelineState` 使用 `sync.Mutex` 保护并发写入 ### 错误处理与降级 | 故障点 | 处理策略 | 客户端表现 | |--------|---------|-----------| -| STT 失败 | 发送 `STT_ERROR`,终止本次请求 | 回退到纯文本模式 | -| LLM 超时(>10s) | 发送 `LLM_TIMEOUT`,取消 TTS | 提示用户重试 | +| STT 失败 | 发送 `STT_ERROR`,Graph 终止 | 回退到纯文本模式 | +| LLM 超时 | 发送 `LLM_TIMEOUT`,取消下游 | 提示用户重试 | | LLM 部分输出后失败 | 已推送的 `llm_chunk` 保留,发送 `error` 通知中断 | 显示已收到的部分文字 | | TTS 失败 | 静默跳过,`llm_done` 正常发送 | 只有文字回复,无语音 | -| interrupt 打断 | cancel context,清空所有流 | 前端清空播放队列 | +| interrupt 打断 | cancel context,Graph 内所有流中断 | 前端清空播放队列 | --- @@ -942,19 +1015,32 @@ type SessionRepository interface { ### 依赖注入 ```go -if cfg.Storage.Driver == "postgres" { - pool, _ := store.NewPostgresPool(ctx, cfg.Storage.DSN) +// 存储层初始化 +if cfg.Storage.Persistence.Enabled { + pool, _ := store.NewPostgresPool(ctx, cfg.Storage.Persistence.DSN) userRepo = store.NewPgUserRepository(pool) msgRepo = store.NewPgMessageRepository(pool) sessRepo = store.NewPgSessionRepository(pool) +} + +// Session Manager 初始化(支持三级存储自动降级) +if cfg.Storage.Redis.Enabled { + sessionMgr = session.NewTieredManager(30*time.Minute, 20, redisClient, + session.WithMessageRepository(msgRepo), + session.WithSessionRepository(sessRepo), + ) +} else if cfg.Storage.Persistence.Enabled { sessionMgr = session.NewMemoryManager(30*time.Minute, 20, session.WithMessageRepository(msgRepo), session.WithSessionRepository(sessRepo), ) } else { - userRepo = store.NewMemUserRepository() sessionMgr = session.NewMemoryManager(30*time.Minute, 20) } + +// Eino Graph 初始化 +pipelineGraph, _ := eino.NewPipelineGraph(ctx, cfg, sttService, ttsService, sessionMgr) +orchestrator := eino.NewEinoOrchestrator(pipelineGraph, sessionMgr, cfg.AI.LLM.Model) ``` --- @@ -1021,27 +1107,27 @@ type AIConfig struct { } type STTConfig struct { - Provider string `mapstructure:"provider"` // "deepgram" | "mimo" | "xiaomi" + Provider string `mapstructure:"provider"` // "mimo" | "deepgram" APIKey string `mapstructure:"api_key"` - Model string `mapstructure:"model"` // 默认 "nova-2" + Model string `mapstructure:"model"` // 默认 "mimo-v2.5-asr" Endpoint string `mapstructure:"endpoint"` Timeout int `mapstructure:"timeout"` // 秒,默认 5 HTTPClientTimeout int `mapstructure:"http_client_timeout"` // 秒,默认 30 } type LLMConfig struct { - Provider string `mapstructure:"provider"` // "openai" + Provider string `mapstructure:"provider"` // "dashscope" / "openai" APIKey string `mapstructure:"api_key"` - Model string `mapstructure:"model"` // 默认 "gpt-4o" + Model string `mapstructure:"model"` // 默认 "qwen3-vl-plus" Endpoint string `mapstructure:"endpoint"` - Timeout int `mapstructure:"timeout"` // 秒,默认 10 + Timeout int `mapstructure:"timeout"` // 秒,默认 30 HTTPClientTimeout int `mapstructure:"http_client_timeout"` // 秒,默认 60 } type TTSConfig struct { - Provider string `mapstructure:"provider"` // "openai" | "mimo" | "xiaomi" + Provider string `mapstructure:"provider"` // "mimo" | "openai" APIKey string `mapstructure:"api_key"` - Model string `mapstructure:"model"` // 默认 "tts-1" + Model string `mapstructure:"model"` // 默认 "mimo-v2.5-tts" Voice string `mapstructure:"voice"` // 默认 "mimo_default" Speed float64 `mapstructure:"speed"` // 默认 1.0 Endpoint string `mapstructure:"endpoint"` @@ -1094,30 +1180,34 @@ redis: ai: stt: - provider: deepgram - model: nova-2 - endpoint: "wss://api.deepgram.com/v1/listen" + provider: mimo + model: mimo-v2.5-asr + endpoint: "https://api.xiaomimimo.com/v1" timeout: 5 http_client_timeout: 30 llm: - provider: openai - model: gpt-4o - endpoint: "https://api.openai.com/v1" - timeout: 10 + provider: dashscope + model: qwen3-vl-plus + endpoint: "https://dashscope.aliyuncs.com/compatible-mode/v1" + timeout: 30 http_client_timeout: 60 tts: - provider: openai - model: tts-1 + provider: mimo + model: mimo-v2.5-tts voice: mimo_default speed: 1.0 - endpoint: "https://api.openai.com/v1" + endpoint: "https://token-plan-cn.xiaomimimo.com/v1" timeout: 5 http_client_timeout: 30 output_format: mp3 sample_rate: 24000 storage: - driver: memory + redis: + enabled: true + persistence: + enabled: true + driver: postgres auth: access_ttl: 15 diff --git a/docs/03-技术选型.md b/docs/03-技术选型.md index ff2060c..fce53c7 100644 --- a/docs/03-技术选型.md +++ b/docs/03-技术选型.md @@ -8,14 +8,16 @@ ``` 技术选型 +├── AI 编排框架 +│ └── CloudWeGo Eino Graph(声明式 DAG 编排,替代手写 goroutine 管道) ├── AI 服务栈 -│ ├── STT: Deepgram(默认) / MiMo ASR -│ ├── LLM: GPT-4o(默认) / 通义千问等 OpenAI 兼容模型 -│ └── TTS: OpenAI TTS(默认) / MiMo TTS +│ ├── STT: MiMo ASR(默认) / Deepgram +│ ├── LLM: DashScope qwen3-vl-plus(默认) / GPT-4o 等 OpenAI 兼容模型 +│ └── TTS: MiMo TTS(默认) / OpenAI TTS ├── 持久化层 │ ├── 数据库: PostgreSQL(pgx/v5,手写 SQL) │ ├── 迁移: 嵌入式 SQL 文件,自动执行 -│ └── 存储模式: Memory(默认)+ Write-Through 到 PG / Redis(可切换) +│ └── 存储模式: 三级存储 TieredManager(L1 Memory → L2 Redis → L3 PostgreSQL) ├── 认证与用户系统 │ ├── 认证方案: JWT (HS256), access 15min + refresh 7day │ ├── JWT 库: golang-jwt/jwt/v5 @@ -30,37 +32,77 @@ --- -## 一、AI 服务栈选型 +## 一、AI 编排框架选型 + +### 候选方案对比 + +| 框架 | 语言 | 特点 | CamTalk 适用性 | +|------|------|------|---------------| +| **CloudWeGo Eino** | Go | Go 原生、类型安全、流式原生、Graph DAG 编排 | ✅ 完美匹配 | +| LangChain Go | Go | 生态丰富但较重,抽象层多 | ❌ 过度抽象 | +| 自研编排 | Go | 完全可控 | ❌ 维护成本高 | + +### 选择 Eino 的理由 + +| 维度 | 手写 goroutine(旧方案) | Eino Graph(新方案) | +|------|------------------------|---------------------| +| 编排方式 | 手动 `go func()` + `sync.WaitGroup` | 声明式 DAG,类型安全 | +| 流式处理 | 自定义 `chan` 传递 | `StreamReader` + `Pipe`,自动转换 | +| 错误处理 | 各节点独立处理,不一致 | Graph 级别统一错误传播 | +| 回调/AOP | 日志散落各处 | `callbacks.Handler` 统一注入 | +| 配置灵活性 | Pipeline 创建时固定 | 每请求 `Option` 动态注入 | +| 可测试性 | 需启动 goroutine | `Graph.Invoke()` 直接测试 | +| 扩展性 | 修改 Pipeline 代码 | 添加节点 + 边,无侵入 | + +### 核心依赖 + +``` +github.com/cloudwego/eino v0.9.9 # 核心框架 +github.com/cloudwego/eino-ext/components/model/openai v0.1.13 # OpenAI 兼容 ChatModel +``` + +**核心理由**: +1. Go 原生,泛型支持,编译时类型检查 +2. 原生流式处理(`StreamReader`),适合 LLM token 级推送 +3. Graph 支持分支、并行、循环,满足当前和未来需求 +4. Callback 机制实现 AOP(日志、指标、消息推送) +5. eino-ext 提供 OpenAI ChatModel 实现,直接对接 DashScope + +> 详细的 Eino 框架使用文档见 [11-Eino框架技术文档](11-Eino框架技术文档.md),重构方案见 [10-Eino重构方案](10-Eino重构方案.md),实施记录见 [12-Eino重构实施记录](12-Eino重构实施记录.md)。 + +--- + +## 二、AI 服务栈选型 ### STT(语音识别) | 方案 | 延迟 | 成本 | 特点 | |------|------|------|------| -| **Deepgram**(默认) | <500ms | 按分钟计费 | 流式识别,延迟极低,WebSocket 接口 | -| **MiMo ASR**(小米) | ~1s | 按量计费 | 国产替代,兼容 OpenAI chat/completions 格式,HTTP 非流式 | +| **MiMo ASR**(默认) | ~1s | 按量计费 | 国产替代,兼容 OpenAI chat/completions 格式,HTTP 非流式 | +| **Deepgram** | <500ms | 按分钟计费 | 流式识别,延迟极低,WebSocket 接口 | | Whisper API | 1-3s | 按分钟计费 | 准确率高,支持多语言 | | FunASR | <500ms | 自部署免费 | 阿里开源,中文优化 | -当前默认使用 Deepgram nova-2,可通过 `ai.stt.provider` 配置切换到 MiMo ASR。 +当前默认使用 MiMo ASR(mimo-v2.5-asr),可通过 `ai.stt.provider` 配置切换到 Deepgram。 ### LLM(多模态大模型) | 方案 | 成本 | 特点 | |------|------|------| -| **GPT-4o**(默认) | $2.5/1M tokens | 视觉理解能力强,API 成熟,流式推理 | -| 通义千问 qwen3-vl-plus | 按量计费 | 阿里云,通过 OpenAI 兼容接口调用 | +| **DashScope qwen3-vl-plus**(默认) | 按量计费 | 阿里云,通过 OpenAI 兼容接口调用,视觉理解能力强 | +| GPT-4o | $2.5/1M tokens | OpenAI,API 成熟,流式推理 | | Claude Sonnet | $3/1M tokens | Anthropic,长上下文能力强 | -代码通过 OpenAI 兼容接口调用,可灵活切换到任何兼容服务商。配置 `ai.llm.provider`、`ai.llm.model`、`ai.llm.endpoint` 即可。 +LLM 通过 Eino 框架的 `eino-ext/components/model/openai` ChatModel 组件接入,支持任何 OpenAI 兼容接口。配置 `ai.llm.provider`、`ai.llm.model`、`ai.llm.endpoint` 即可切换。 ### TTS(语音合成) | 方案 | 成本 | 特点 | |------|------|------| -| **OpenAI TTS**(默认) | $15/1M 字符 | 音质自然,支持流式,默认模型 tts-1,语音 alloy | -| MiMo TTS(小米) | 按量计费 | 国产替代,通过配置切换 | +| **MiMo TTS**(默认) | 按量计费 | 国产替代,通过配置切换,模型 mimo-v2.5-tts | +| OpenAI TTS | $15/1M 字符 | 音质自然,支持流式,默认模型 tts-1,语音 alloy | -当前默认使用 OpenAI TTS(tts-1, alloy),可通过 `ai.tts.provider` 配置切换。 +当前默认使用 MiMo TTS(mimo-v2.5-tts),可通过 `ai.tts.provider` 配置切换到 OpenAI TTS。 --- @@ -149,17 +191,19 @@ ORDER BY created_at DESC LIMIT 20; ``` -### 冷热分离架构 +### 冷热分离架构(三级存储) ``` -Go Gateway - ├── 写入路径 → Redis(实时会话状态) - │ → PostgreSQL(对话历史 + 用量) - └── 读取路径 → Redis(当前上下文,快) - → PostgreSQL(历史记录,慢) +Go Gateway (TieredManager) + ├── L1: Memory(进程内缓存,微秒级) + ├── L2: Redis(分布式缓存,毫秒级) + └── L3: PostgreSQL(持久化存储,冷数据) + +读取路径:L1 → L2 → L3,逐级回源,命中后向上回填 +写入路径:L1 → L2(同步) → L3(异步) ``` -建议异步写入——实时消息先写 Redis(快),异步批量刷入 PostgreSQL(慢),不影响对话体验。 +`TieredManager` 自动管理三级存储,后台 goroutine 每 30 秒 ping Redis 健康状态,Redis 故障时自动降级为 L1+L3 模式。 ### 决策流程 diff --git a/docs/05-语音交互.md b/docs/05-语音交互.md index fcbb1aa..4e88829 100644 --- a/docs/05-语音交互.md +++ b/docs/05-语音交互.md @@ -14,6 +14,8 @@ 麦克风 → VAD → STT → LLM → TTS → 扬声器 ``` +> 后端 AI 编排基于 Eino Graph 声明式 DAG 实现:`START → STT → History → ChatModel → Msg2Str → Splitter → TTS → Done → END`。详见 [11-Eino框架技术文档](11-Eino框架技术文档.md)。 + ## 环节一:VAD(语音活动检测) 从持续音频流中检测"人什么时候在说话",避免将环境噪音当作有效输入。**浏览器端完成**,节省 ~70% 带宽。 @@ -41,12 +43,12 @@ vad.start(); | 方案 | 延迟 | 成本 | 特点 | |------|------|------|------| +| **MiMo ASR**(默认) | ~1s | 按量计费 | 国产替代,兼容 OpenAI 格式,HTTP 非流式 | +| **Deepgram** | <500ms | 按分钟计费 | 流式识别,延迟极低 | | Whisper API | 1-3s | 按分钟计费 | 准确率高,支持多语言 | -| **Deepgram**(默认) | <500ms | 按分钟计费 | 流式识别,延迟极低 | -| **MiMo ASR**(小米) | ~1s | 按量计费 | 国产替代,兼容 OpenAI 格式,HTTP 非流式 | | 浏览器原生 | ~1s | 免费 | 中文效果一般 | -当前实现为**一次性语音识别**(非流式):前端 VAD 检测到用户说完后,将完整音频片段发送到后端,后端调用 `stt.Recognize()` 一次性返回识别结果。流式 STT 为未来优化方向。 +当前实现为**一次性语音识别**(非流式):前端 VAD 检测到用户说完后,将完整音频片段发送到后端,后端通过 Eino Graph 的 STT Lambda 节点调用 `stt.Recognize()` 一次性返回识别结果。流式 STT 为未来优化方向。 音频编码格式:前端 `audio.ts` 将 Float32Array 转为 Int16 PCM(16kHz, pcm_s16le)再编码为 Base64。 @@ -59,8 +61,8 @@ vad.start(); 当前实现参数:Voice `"mimo_default"`(可通过配置切换)、Speed `1.0`、OutputFmt `"mp3"`、SampleRate `24000`。 方案选择: -- **OpenAI TTS**(默认):音质好,延迟中等,按字符计费,模型 tts-1 -- **MiMo TTS**(小米):国产替代,通过配置切换 +- **MiMo TTS**(默认):国产替代,模型 mimo-v2.5-tts,通过配置切换 +- **OpenAI TTS**:音质好,延迟中等,按字符计费,模型 tts-1 - **Edge TTS**(待实现):微软免费方案,音质不错,延迟略高 ## 延迟优化要点 diff --git a/docs/07-成本控制.md b/docs/07-成本控制.md index 7f28405..5e2213d 100644 --- a/docs/07-成本控制.md +++ b/docs/07-成本控制.md @@ -50,12 +50,12 @@ const ACTIVE_INTERVAL = 1000; // 用户说话时 1 秒一帧 ``` 用户提问 → 问题复杂度判断 - ├── 简单识别 → GPT-4o-mini ($0.15/1M tokens) - ├── 深度分析 → GPT-4o ($2.5/1M tokens) - └── 代码/推理 → o1 ($15/1M tokens) + ├── 简单识别 → 轻量模型(如 qwen-turbo) + ├── 深度分析 → qwen3-vl-plus(默认,按量计费) + └── 代码/推理 → 更强模型(如 o1) ``` -> 当前 MVP 阶段使用单一模型(默认 GPT-4o),模型分级路由为未来优化方向。通过配置 `ai.llm.model` 可手动切换模型。 +> 当前 MVP 阶段使用单一模型(默认 DashScope qwen3-vl-plus),模型分级路由为未来优化方向。通过配置 `ai.llm.model` 可手动切换模型。LLM 通过 Eino 框架的 eino-ext ChatModel 组件接入,支持任何 OpenAI 兼容接口。 ## 策略四:缓存与复用(待实现) diff --git a/docs/09-技术名词解释.md b/docs/09-技术名词解释.md index 663baaa..97732c7 100644 --- a/docs/09-技术名词解释.md +++ b/docs/09-技术名词解释.md @@ -34,3 +34,16 @@ | **STT** | 语音转文字 | Speech-to-Text。Deepgram 流式识别延迟 <500ms。备选 FunASR(阿里开源,可自部署)。 | | **TTS** | 文字转语音 | Text-to-Speech。OpenAI TTS 音质接近真人。Edge TTS 免费。支持流式——边生成边读,不必等全部生成完。 | | **GPT-4o-mini** | 轻量分类模型 | 又快又便宜的小模型,用于模型路由——先用小模型判断问题复杂度,简单问题走小模型省 API 费用。 | + +## AI 编排框架相关 + +| 名词 | 一句话 | 展开 | +|------|--------|------| +| **Eino** | 字节跳动开源的 Go AI 应用开发框架 | CloudWeGo Eino,提供 Graph DAG 编排、组件抽象(ChatModel/Tool 等)、流式处理(StreamReader)和 Callback AOP 机制。CamTalk 用它替代手写 goroutine 管道。 | +| **compose.Graph** | Eino 的 DAG 编排器 | 声明式有向无环图,节点可以是 Lambda、ChatModel、ToolsNode 等,边定义数据流向。支持分支(AddBranch)、并行和循环。 | +| **Lambda** | Graph 中的可组合函数单元 | 四种模式:InvokableLambda(同步)、StreamableLambda(流式输出)、CollectableLambda(流式输入)、TransformableLambda(双向流式)。 | +| **StreamReader** | Eino 的流式数据抽象 | `schema.StreamReader[T]`,类似 io.Reader 的语义,`Recv()` 读取一帧,`io.EOF` 表示流结束。`schema.Pipe[T]()` 创建 StreamReader + StreamWriter 对。 | +| **Callback** | Eino 的 AOP 机制 | 类似中间件的钩子,支持节点生命周期回调(OnStart/OnEnd/OnError/OnEndWithStreamOutput)。CamTalk 用它实现 LLM token 实时推送到客户端。 | +| **ChatModel** | Eino 的 LLM 组件抽象 | 统一接口 `Generate()` 和 `Stream()`,eino-ext 提供 OpenAI 兼容实现,通过 BaseURL 可对接 DashScope 等兼容接口。 | +| **eino-ext** | Eino 的组件扩展库 | 提供具体组件实现:OpenAI ChatModel、各种 Tool Backend 等。CamTalk 使用 `eino-ext/components/model/openai`。 | +| **PipelineState** | Graph 级别的共享状态 | 通过 `compose.WithGenLocalState` 注册,每请求独立实例,线程安全(sync.Mutex),跨节点共享数据(如 LLM 完整回复、Token 用量)。 | diff --git a/docs/10-Eino重构方案.md b/docs/10-Eino重构方案.md index f7f9a3e..a00b15b 100644 --- a/docs/10-Eino重构方案.md +++ b/docs/10-Eino重构方案.md @@ -1,7 +1,7 @@ # CamTalk 后端 AI 编排层 Eino 重构方案 > 创建日期:2026-06-19 -> 状态:草案 +> 状态:已实施(实施记录见 [12-Eino重构实施记录](12-Eino重构实施记录.md)) ## 1. 背景与目标 diff --git a/docs/12-Eino重构实施记录.md b/docs/12-Eino重构实施记录.md deleted file mode 100644 index cafcdca..0000000 --- a/docs/12-Eino重构实施记录.md +++ /dev/null @@ -1,204 +0,0 @@ -# CamTalk Eino 重构实施记录 - -> 创建日期:2026-06-19 -> 状态:已完成 - -## 1. 重构背景 - -CamTalk 原 AI 编排层(`internal/orchestrator/pipeline.go`)使用手写 goroutine + WaitGroup + channel 实现 STT → LLM → TTS 流式管道,存在以下问题: - -1. **编排逻辑硬编码**:流程写死在 `ProcessQuery()` 中,扩展需重写 goroutine 调度 -2. **并发控制粗糙**:手动 `go func()` + `sync.WaitGroup`,缺乏结构化流式传递 -3. **无回调/AOP 机制**:日志、指标、追踪散落各处 -4. **配置耦合**:模型名、TTS 参数硬编码在 Pipeline 结构体 -5. **错误处理不一致**:TTS 错误被静默吞掉,缺乏统一模式 - -**重构目标**:使用 Eino Graph 替换手写 Pipeline,实现声明式编排、统一回调、按请求动态配置,保持 WebSocket 协议和 REST API 不变。 - -## 2. 整体架构变更 - -### 2.1 重构前 - -``` -WS Handler → Orchestrator.Pipeline.ProcessQuery() - ├→ goroutine: STT.Recognize() - ├→ goroutine: LLM.ChatStream() ──→ chan chunk ──→ Sender - └→ goroutine: Splitter → TTS.SynthesizeStream() ──→ chan audio ──→ Sender - WaitGroup.Wait() - Sender.SendLLMDone() -``` - -### 2.2 重构后 - -``` -WS Handler → EinoOrchestrator.ProcessQuery() - ├→ Graph.Stream(ctx, input) - │ ├→ STT Lambda ─→ History Lambda ─→ ChatModel ─→ Splitter ─→ TTS ─→ Done - │ │ (State 写入) (Callback (Transform) (Invoke) (Invoke) - │ │ 流式推送) - │ └→ 消费 StreamReader(触发整条链路惰性执行) - └→ 追加助手消息到历史 -``` - -### 2.3 关键设计决策 - -| 决策 | 选择 | 理由 | -|------|------|------| -| Graph 调用模式 | Stream | ChatModel 需要真正的 token 级流式输出 | -| LLM 组件 | eino-ext ChatModel | 原生 Eino 组件,直接对接 DashScope | -| 消息推送 | Callback(LLM)+ Sender(其他) | LLM token 流式推送需要 Callback | -| 值类型 vs 指针 | 值类型统一 | 避免框架类型转换不匹配 | -| 历史追加 | 适配器负责 | Done 节点只负责发送 llm_done | - -## 3. 分阶段实施 - -### Phase 1:基础设施(提交 `fd5c771`) - -**目标**:引入 Eino 依赖,创建基础类型和 Callback。 - -**任务清单**: - -| 任务 | 文件 | 说明 | -|------|------|------| -| 引入 Eino 依赖 | `go.mod` | `eino v0.9.9` + `eino-ext/components/model/openai v0.1.13` | -| 数据类型定义 | `eino/types.go` | `PipelineInput`、`PipelineOutput`、`STTOutput`、`TokenUsage` | -| State 定义 | `eino/state.go` | `PipelineState` 含 `sync.Mutex` 并发保护 | -| 消息推送 Callback | `eino/callback.go` | `BuildCallbackHandler()` 使用 `callbacks.NewHandlerHelper()` | - -**关键实现**: -- `PipelineState` 使用 `strings.Builder` + `sync.Mutex` 累积 LLM 完整回复 -- Callback 通过 `ModelCallbackHandler.OnEndWithStreamOutput` 逐 token 推送 `llm_chunk` -- Sender/RequestID/PipelineState 通过 `context.WithValue` 注入 - -**验证**:`go build ./cmd/server` ✓ - ---- - -### Phase 2:节点实现(提交 `fd5c771`) - -**目标**:实现 Graph 中的 5 个 Lambda 节点。 - -**任务清单**: - -| 任务 | 文件 | Lambda 类型 | 输入 → 输出 | -|------|------|------------|------------| -| STT Lambda | `eino/nodes_stt.go` | InvokableLambda | `PipelineInput → STTOutput` | -| 历史组装 Lambda | `eino/nodes_history.go` | InvokableLambda | `STTOutput → []*schema.Message` | -| 句子分割 Lambda | `eino/nodes_splitter.go` | TransformableLambda | `StreamReader[string] → StreamReader[[]string]` | -| TTS Lambda | `eino/nodes_tts.go` | InvokableLambda | `[]string → struct{}` | -| Done Lambda | `eino/nodes_done.go` | InvokableLambda | `struct{} → PipelineOutput` | - -**关键实现**: -- STT 节点将输入元数据写入 State,供下游节点读取 -- History 节点从 State 读取 SessionID/Scenario/ImageData,构建系统提示词 + 多模态消息 -- Splitter 使用 `TransformableLambda` 按句子分隔符切分,逐句输出给 TTS -- TTS 节点调用 `ttsService.SynthesizeStream()`,逐 chunk 推送 `tts_audio` -- Done 节点从 State 读取完整回复,发送 `llm_done` -- 所有 Lambda 使用值类型(非指针),返回 `*compose.Lambda` - -**验证**:`go build ./internal/eino/...` ✓ - ---- - -### Phase 3:Graph 构建与适配器(提交 `4b731b5`) - -**目标**:构建 Graph、实现适配器、切换 main.go。 - -**任务清单**: - -| 任务 | 文件 | 说明 | -|------|------|------| -| Graph 构建 | `eino/graph.go` | `NewPipelineGraph()` 组装 6 个节点 + 边 + 编译 | -| 适配器 | `eino/adapter.go` | `EinoOrchestrator` 实现 `orchestrator.Orchestrator` 接口 | -| main.go 切换 | `cmd/server/main.go` | 移除旧 LLM + orchestrator,替换为 Eino | - -**Graph 拓扑**: -``` -START → STT → History → ChatModel → Splitter → TTS → Done → END -``` - -**适配器职责**: -1. 解码 base64 音频/图片 -2. 获取会话配置 -3. 注入 Sender/RequestID/SessionID/StartTime/State 到 context -4. 追加用户消息到历史 -5. 调用 `graph.Stream(ctx, input, callbacks)` 触发惰性执行 -6. 消费 `StreamReader[PipelineOutput]` -7. 追加助手消息到历史 - -**关键实现**: -- eino-ext ChatModel 配置:`BaseURL` 对接 DashScope,`Timeout` 控制请求超时 -- Callback 在运行时通过 `compose.WithCallbacks()` 传入,不在编译时注册 -- 元数据(SessionID/Scenario 等)通过 State 跨节点传递,不通过 Graph 边传递 - -**变更文件**: -- 修改 `state.go`:新增 SessionID/RequestID/ImageData 等字段 -- 修改 `nodes_stt.go`:写入元数据到 State -- 修改 `nodes_history.go`:从 State 读取元数据(移除 HistoryInput 依赖) -- 修改 `nodes_done.go`:移除历史追加(由适配器负责) - -**验证**:`go build ./cmd/server` ✓,`go vet ./...` ✓ - ---- - -### Phase 4:清理与测试(提交 `4ffd845`) - -**目标**:删除旧代码,编写单元测试。 - -**删除的文件**: - -| 文件 | 说明 | -|------|------| -| `orchestrator/pipeline.go` | 旧 STT→LLM→TTS 手写 goroutine 管道(-547 行) | -| `orchestrator/splitter.go` | 旧句子切分器(-114 行) | -| `orchestrator/pipeline_test.go` | 旧 Pipeline 测试(-309 行) | -| `ai/llm/openai.go` | 旧 LLM OpenAI 实现(-548 行) | -| `ai/llm/openai_test.go` | 旧 LLM 测试(-143 行) | - -**保留的文件**: - -| 文件 | 保留原因 | -|------|---------| -| `orchestrator/orchestrator.go` | Orchestrator 接口(ws/handler 依赖) | -| `orchestrator/sender.go` | Sender 接口(eino/callback 依赖) | -| `ai/llm/llm.go` | Request/Chunk/TokenUsage 类型定义 | -| `ai/llm/prompt.go` | BuildSystemPrompt(eino/nodes_history 依赖) | -| `ai/llm/scenarios.go` | GetScenarioPrompt(eino/nodes_history 依赖) | - -**新增测试**:`eino/graph_test.go`(13 个测试) - -| 测试 | 覆盖内容 | -|------|---------| -| `TestDetectImageMimeType` | JPEG/PNG/GIF/WebP/未知格式检测 | -| `TestBuildPipelineInput` | 文本输入构建 | -| `TestBuildPipelineInput_WithAudioData` | 音频+图片输入构建 | -| `TestPipelineState_AppendAndGet` | State 文本追加和读取 | -| `TestPipelineState_ConcurrentAccess` | State 并发安全(100 goroutine) | -| `TestContextInjection` | Sender/RequestID/State 注入和提取 | -| `TestLatencyFromCtx` | 延迟计算 | -| `TestEinoOrchestrator_ImplementsInterface` | 接口实现检查 | -| `TestNew*Lambda_ReturnsNonNil` | 5 个 Lambda 构造函数非空检查 | - -**验证**:`go build ./...` ✓,`go vet ./...` ✓,`go test ./...` ✓ - -## 4. 代码变更统计 - -| 阶段 | 提交 | 新增 | 删除 | 净变化 | -|------|------|------|------|--------| -| Phase 1 + 2 | `fd5c771` | +946 | -24 | +922 | -| Phase 3 | `4b731b5` | +395 | -98 | +297 | -| Phase 4 | `4ffd845` | +235 | -1661 | -1426 | -| **合计** | | **+1576** | **-1783** | **-207** | - -重构后代码量净减少 207 行,同时获得了更好的可维护性、可测试性和可扩展性。 - -## 5. 遗留事项 - -| 事项 | 优先级 | 说明 | -|------|--------|------| -| eino-ext ChatModel DashScope 兼容性端到端验证 | 高 | 需要真实 API Key 验证流式输出和多模态 | -| LLM 超时控制 | 中 | eino-ext ChatModel 的 `Timeout` 配置需验证 | -| TTS 流式优化 | 中 | 当前 TTS 是 InvokableLambda,可改为 StreamableLambda | -| ReAct Agent 扩展 | 低 | 基于 Graph Branch 实现工具调用循环 | -| Model Router | 低 | 按场景/成本路由不同 LLM | -| 指标监控 | 低 | 通过 Callback 接入 Prometheus | diff --git a/docs/README.md b/docs/README.md index 3f8ad60..a898a9a 100644 --- a/docs/README.md +++ b/docs/README.md @@ -7,20 +7,25 @@ CamTalk 是一款多模态实时 AI 视觉对话助手。用户通过摄像头 | 文档 | 说明 | |------|------| | [01-架构设计](01-架构设计.md) | 系统架构、技术栈、模块设计、数据库、部署架构(含 Mermaid 图) | -| [02-接口文档](02-接口文档.md) | WebSocket 协议、REST API、AI 服务层、编排器、Session Manager、配置管理、数据模型、错误码 | -| [03-技术选型](03-技术选型.md) | 各技术的选型对比与决策理由 | +| [02-接口文档](02-接口文档.md) | WebSocket 协议、REST API、AI 服务层、Eino 编排器、Session Manager、配置管理、数据模型、错误码 | +| [03-技术选型](03-技术选型.md) | 各技术的选型对比与决策理由(含 Eino 框架选型) | | [04-用户故事](04-用户故事.md) | P0/P1/P2 用户故事、验收标准 | | [05-语音交互](05-语音交互.md) | VAD → STT → LLM → TTS 全链路、延迟优化 | | [06-视觉理解](06-视觉理解.md) | 帧采样策略、图像编码、多模态 LLM 输入机制 | | [07-成本控制](07-成本控制.md) | 智能采样、端云协同、模型分级、缓存复用 | | [08-功能创意](08-功能创意.md) | 未来功能创意清单 | -| [09-技术名词解释](09-技术名词解释.md) | 前端/后端/AI 服务技术名词简明解释 | +| [09-技术名词解释](09-技术名词解释.md) | 前端/后端/AI 服务/Eino 框架技术名词简明解释 | +| [10-Eino重构方案](10-Eino重构方案.md) | Eino Graph 替换手写 goroutine 管道的设计方案 | +| [11-Eino框架技术文档](11-Eino框架技术文档.md) | Eino 框架在 CamTalk 中的使用指南(Graph、Lambda、Callback、State) | + ## 推荐阅读顺序 1. **01-架构设计** — 理解三层架构、技术栈和模块全貌 2. **02-接口文档** — 前后端通信契约,实现时的最高依据 -3. **03-技术选型** — 了解为什么选这些技术 +3. **03-技术选型** — 了解为什么选这些技术(含 Eino 框架) 4. **04-用户故事** — 明确功能优先级 5. **05~07** — 各技术领域的详细设计 6. **09-技术名词解释** — 遇到不熟悉的名词时查阅 +7. **10~12** — Eino 重构相关(方案、框架文档、实施记录) + diff --git a/frontend/src/App.css b/frontend/src/App.css index e31c238..7d8bdd0 100644 --- a/frontend/src/App.css +++ b/frontend/src/App.css @@ -969,6 +969,7 @@ body { .chat-input__field { flex: 1; + min-width: 0; padding: 10px 14px; border: 1px solid var(--color-border); border-radius: var(--radius-sm); @@ -996,6 +997,7 @@ body { font-size: 0.9rem; cursor: pointer; transition: opacity var(--transition-fast); + flex-shrink: 0; } .chat-input__send:hover:not(:disabled) { @@ -1503,6 +1505,48 @@ body { color: white; } +/* ---- Text-only mode (video ended, chat preserved) ---- */ + +.video-controls__text-only { + display: flex; + flex-direction: column; + align-items: center; + gap: 10px; + width: 100%; +} + +.video-ended-hint { + display: flex; + flex-direction: column; + align-items: center; + gap: 2px; + padding: 4px 0; +} + +.video-ended-hint__title { + font-size: 0.85rem; + font-weight: 600; + color: var(--color-text-muted); +} + +.video-ended-hint__sub { + font-size: 0.7rem; + color: var(--color-text-muted); + opacity: 0.7; +} + +.btn--outline { + background: transparent; + border: 1px solid var(--color-border); + color: var(--color-text-muted); +} + +.btn--outline:hover { + background: var(--color-surface-2); + color: var(--color-text); + border-color: var(--color-surface-3); +} + /* ---- Scene Cards (empty state) ---- */ .scene-cards { diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index 389b569..0c3d1d8 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -80,6 +80,7 @@ function AppContent() { isObserving, startSession, stopSession, + stopVideo, interrupt, isCameraOn, isMicOn, @@ -397,9 +398,9 @@ function AppContent() { - ) : ( + ) : isCameraOn ? ( <> - {/* 通话态:核心控制工具栏 */} + {/* 视频通话态:核心控制工具栏 */}
{/* 通话态模式切换 */} @@ -456,6 +457,24 @@ function AppContent() { + ) : ( + <> + {/* 文字对话态:视频已结束 */} +
+
+ 📹 {tr("video.ended")} + {tr("video.ended.hint")} +
+
+ + +
+
+ )} @@ -501,6 +520,7 @@ function AppContent() { isProcessing={isProcessing} isVADReady={isVADReady} vadError={vadError ?? undefined} + isCameraOn={isCameraOn} isMicOn={isMicOn} isSpeaking={isSpeaking} onSendText={sendTextMessage} diff --git a/frontend/src/components/ChatPanel/index.tsx b/frontend/src/components/ChatPanel/index.tsx index 41a945c..4336803 100644 --- a/frontend/src/components/ChatPanel/index.tsx +++ b/frontend/src/components/ChatPanel/index.tsx @@ -18,6 +18,7 @@ interface ChatPanelProps { isProcessing?: boolean; isVADReady?: boolean; vadError?: string | null; + isCameraOn?: boolean; isMicOn?: boolean; isSpeaking?: boolean; onSendText?: (text: string) => void; @@ -44,6 +45,7 @@ export function ChatPanel({ isProcessing, isVADReady, vadError, + isCameraOn, isMicOn, isSpeaking, onSendText, @@ -198,7 +200,7 @@ export function ChatPanel({ )} {/* VAD 初始化中 */} - {isConnected && isVADReady === false && !vadError && ( + {isConnected && isCameraOn && isVADReady === false && !vadError && (
diff --git a/frontend/src/hooks/useVisionSession.ts b/frontend/src/hooks/useVisionSession.ts index 6b474a1..0e56558 100644 --- a/frontend/src/hooks/useVisionSession.ts +++ b/frontend/src/hooks/useVisionSession.ts @@ -363,6 +363,22 @@ export function useVisionSession(accessToken?: string | null) { setIsMicOn(false); }, [stopObserving, stopVAD, stopMic, stopCamera, disconnect]); + /** 结束视频,保留聊天和连接 */ + const stopVideo = useCallback(async () => { + stopObserving(); + setMode("dialogue"); + await stopVAD(); + stopMic(); + stopCamera(); + ttsPlayerRef.current?.stop(); + setIsAudioPlaying(false); + setCurrentReply(""); + setIsProcessing(false); + setIsCameraOn(false); + setIsMicOn(false); + // 不断开 WebSocket,不清空消息、历史、统计 + }, [stopObserving, stopVAD, stopMic, stopCamera]); + /** 摄像头开关 */ const toggleCamera = useCallback(async () => { if (isCameraOn) { @@ -422,12 +438,7 @@ export function useVisionSession(accessToken?: string | null) { // 未连接时:自动连接,消息加入待发队列 if (statusRef.current !== "connected") { pendingMessagesRef.current.push({ text: text.trim(), requestId }); - // 添加用户消息到 UI(立即反馈) - setMessages((prev) => [ - ...prev, - { role: "user", content: text.trim(), timestamp: Date.now() }], - ); - // 自动连接 WebSocket + // 自动连接 WebSocket(消息在连接成功后由 flush 统一添加到 UI,避免重复) connect(accessToken || undefined); return; } @@ -479,6 +490,7 @@ export function useVisionSession(accessToken?: string | null) { toggleMode, startSession, stopSession, + stopVideo, interrupt, isCameraOn, isMicOn, diff --git a/frontend/src/lib/i18n/en-US.ts b/frontend/src/lib/i18n/en-US.ts index 9d4087d..ed2e92c 100644 --- a/frontend/src/lib/i18n/en-US.ts +++ b/frontend/src/lib/i18n/en-US.ts @@ -51,6 +51,8 @@ export const enUS: TranslationMap = { "video.placeholder": "Type in the chat panel to start", "video.cameraOff": "Camera is off", "video.cameraOff.hint": "You can type in the chat panel", + "video.ended": "Video Ended", + "video.ended.hint": "You can continue chatting below", // Controls "controls.connecting": "Connecting...", @@ -64,6 +66,9 @@ export const enUS: TranslationMap = { "controls.observation": "👁️ Observe", "controls.interrupt": "⏹ Interrupt", "controls.stop": "End Session", + "controls.stopVideo": "End Video", + "controls.endSession": "End Session", + "controls.resumeVideo": "📹 Resume Video", // Chat panel "chat.title": "Chat", diff --git a/frontend/src/lib/i18n/ja-JP.ts b/frontend/src/lib/i18n/ja-JP.ts index 574fdb3..5d6e1db 100644 --- a/frontend/src/lib/i18n/ja-JP.ts +++ b/frontend/src/lib/i18n/ja-JP.ts @@ -51,6 +51,8 @@ export const jaJP: TranslationMap = { "video.placeholder": "右側のチャットに入力して開始", "video.cameraOff": "カメラがオフです", "video.cameraOff.hint": "右側のチャットでテキスト対話ができます", + "video.ended": "ビデオ終了", + "video.ended.hint": "下にテキストを入力して会話を続けることができます", // Controls "controls.connecting": "接続中...", @@ -64,6 +66,9 @@ export const jaJP: TranslationMap = { "controls.observation": "👁️ 観察モード", "controls.interrupt": "⏹ 中断", "controls.stop": "対話を終了", + "controls.stopVideo": "ビデオ終了", + "controls.endSession": "セッション終了", + "controls.resumeVideo": "📹 ビデオ再開", // Chat panel "chat.title": "チャット", diff --git a/frontend/src/lib/i18n/zh-CN.ts b/frontend/src/lib/i18n/zh-CN.ts index ca939ab..8d7d686 100644 --- a/frontend/src/lib/i18n/zh-CN.ts +++ b/frontend/src/lib/i18n/zh-CN.ts @@ -51,6 +51,8 @@ export const zhCN: TranslationMap = { "video.placeholder": "在右侧聊天框输入即可开始对话", "video.cameraOff": "摄像头未开启", "video.cameraOff.hint": "可在右侧聊天框打字对话", + "video.ended": "视频已结束", + "video.ended.hint": "您可以继续在下方输入文字对话", // Controls "controls.connecting": "连接中...", @@ -64,6 +66,9 @@ export const zhCN: TranslationMap = { "controls.observation": "👁️ 观察模式", "controls.interrupt": "⏹ 打断", "controls.stop": "结束对话", + "controls.stopVideo": "结束视频", + "controls.endSession": "结束会话", + "controls.resumeVideo": "📹 重新开始视频", // Chat panel "chat.title": "对话",