feat: 对话场景优化 #138

Merged
cfy777 merged 6 commits from fea/models into develop 2026-06-20 14:49:22 +08:00
9 changed files with 348 additions and 150 deletions
Showing only changes of commit 3dc2015a91 - Show all commits

View File

@@ -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-4oLLM、DeepgramSTT、OpenAI TTS。仅通过 Go 网关访问,浏览器不直连。
2. **Go 网关**Gin, gorilla/websocket, Viper, Zap—— WebSocket 服务器、会话管理、AI 编排(基于 CloudWeGo Eino Graph。每个 WebSocket 连接一个 goroutine。
3. **云端 AI 服务** —— 通过 OpenAI 兼容接口可灵活切换。默认:DashScope qwen3-vl-plusLLM、MiMo ASRSTT、MiMo TTSTTS。仅通过 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、SessionRepositoryPostgreSQL + 内存实现。
## 技术栈
@@ -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=<access_token>&conversation_id=<uuid>`
所有消息为 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/Redis30 分钟 TTL |
| AI Orchestrator | STT→LLM→TTS 流式并行管道编排 |
| AI Service Layer | AI 服务抽象层STT/LLM/TTS 多 provider |
| REST API | 健康检查、会话管理Gin 路由 |
| WebSocket Handler | 连接管理、JWT 认证、单播消息推送 |
| Session Manager | 会话状态、对话历史(三级存储:Memory/Redis/PostgreSQL30 分钟 TTL |
| Eino 编排层 | 基于 Eino Graph 的声明式 AI 编排7 节点 DAGStream 模式Callback AOP |
| AI Orchestrator | `EinoOrchestrator` 适配器,包装 Graph 实现 `Orchestrator` 接口 |
| AI Service Layer | AI 服务抽象层STT/TTS 多 providerLLM 通过 eino-ext ChatModel |
| Auth | JWT 双 token 轮转认证bcrypt 密码哈希 |
| Store | 持久化存储层UserRepository/MessageRepository/SessionRepository内存 + PostgreSQL |
| REST API | 健康检查、认证、对话管理Gin 路由) |
| Models | 数据模型定义 |
| Migrations | 数据库版本化迁移(嵌入式 SQL |
| Model Router | 按请求选择 AI 模型(规划中) |
| Rate Limiter | 按用户的令牌桶速率限制(规划中) |

View File

@@ -27,7 +27,7 @@ graph TB
subgraph Gateway["Go 网关"]
WS["WebSocket Handler<br/>连接管理 / 消息分发"]
Session["Session Manager<br/>会话状态 / 对话历史"]
Orch["AI Orchestrator<br/>STT→LLM→TTS 流式并行"]
Orch["AI Orchestrator<br/>Eino Graph 声明式编排"]
Auth["Auth 模块<br/>JWT / bcrypt"]
REST["REST API<br/>健康检查 / 对话管理"]
Store["Store 层<br/>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 LLMChatModel
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: 句子文本
G->>G: Msg2Str + Splitter Lambda句子切分
G->>T: TTS Lambda逐句合成
T-->>G: 音频 chunk
G-->>B: tts_audio {audio}
end
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<br/>会话生命周期"]
ORCH["Orchestrator<br/>STT→LLM→TTS 编排"]
ORCH["EinoOrchestrator<br/>Eino Graph 编排"]
AS["Auth Service<br/>注册/登录/刷新/登出"]
end
subgraph Eino_Layer["Eino 编排层"]
PG["PipelineGraph<br/>7 节点 DAG"]
CB["Callback Handler<br/>LLM token 推送"]
ST["PipelineState<br/>跨节点状态"]
end
subgraph AI_Layer["AI 服务层"]
STT_S["STT Service<br/>Deepgram / MiMo"]
LLM_S["LLM Service<br/>OpenAI 兼容"]
TTS_S["TTS Service<br/>OpenAI / MiMo"]
STT_S["STT Service<br/>MiMo / Deepgram"]
LLM_S["ChatModel<br/>eino-ext OpenAI 兼容"]
TTS_S["TTS Service<br/>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 分钟 TTLWrite-Through 到 PG |
| AI Orchestrator | 编排 STT→LLM→TTS 流式并行管道context 取消 + 超时控制 + 句子切分 |
| Eino 编排层 | 基于 CloudWeGo Eino Graph 的声明式 AI 编排。7 节点 DAGSTT→History→ChatModel→Msg2Str→Splitter→TTS→DoneStream 模式调用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 双实现 |
@@ -306,10 +321,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 存"冷数据"(历史记录)
## 认证设计

View File

@@ -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 流式输出逐 tokenTTS 需要完整句子才能合成。解法:**句子切分器 + 管道并行**
AI 编排层基于 [CloudWeGo Eino](https://github.com/cloudwego/eino) 框架的 `compose.Graph` 实现,替代了原有的手写 goroutine 管道。Eino Graph 是一个声明式的有向无环图DAG编排器支持类型安全的流式数据传递和 Callback AOP 机制
**核心矛盾**LLM 流式输出逐 tokenTTS 需要完整句子才能合成。解法:**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 contextGraph 内所有流中断 | 前端清空播放队列 |
---
@@ -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

View File

@@ -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
├── 持久化层
│ ├── 数据库: PostgreSQLpgx/v5手写 SQL
│ ├── 迁移: 嵌入式 SQL 文件,自动执行
│ └── 存储模式: Memory默认+ Write-Through 到 PG / Redis可切换
│ └── 存储模式: 三级存储 TieredManagerL1 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 ASRmimo-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 | OpenAIAPI 成熟,流式推理 |
| 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 TTStts-1, alloy),可通过 `ai.tts.provider` 配置切换。
当前默认使用 MiMo TTSmimo-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 模式
### 决策流程

View File

@@ -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 PCM16kHz, 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**(待实现):微软免费方案,音质不错,延迟略高
## 延迟优化要点

View File

@@ -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 兼容接口。
## 策略四:缓存与复用(待实现)

View File

@@ -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 用量)。 |

View File

@@ -1,7 +1,7 @@
# CamTalk 后端 AI 编排层 Eino 重构方案
> 创建日期2026-06-19
> 状态:草案
> 状态:已实施(实施记录见 [12-Eino重构实施记录](12-Eino重构实施记录.md)
## 1. 背景与目标

View File

@@ -7,20 +7,24 @@ 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 |
| [12-Eino重构实施记录](12-Eino重构实施记录.md) | Eino 重构的四阶段实施记录、代码变更统计、测试覆盖 |
## 推荐阅读顺序
1. **01-架构设计** — 理解三层架构、技术栈和模块全貌
2. **02-接口文档** — 前后端通信契约,实现时的最高依据
3. **03-技术选型** — 了解为什么选这些技术
3. **03-技术选型** — 了解为什么选这些技术(含 Eino 框架)
4. **04-用户故事** — 明确功能优先级
5. **05~07** — 各技术领域的详细设计
6. **09-技术名词解释** — 遇到不熟悉的名词时查阅
7. **10~12** — Eino 重构相关(方案、框架文档、实施记录)