topfans/docs/specs/2026-06-29-ai-chat-dify-mvp-design.md
Lenticular Studio Agent 65ce6bba12 feat: Dify 部署脚本修复 + AI 搭子 MVP 接入
主要改动:

fix(docker/dify-deploy): 修复脚本核心功能
- heredoc 单引号 bug: 'ENVEOF' 改为 ENVEOF,变量正确展开
- 端口默认值 8083/8084/8085 对齐 .env.prod 生产配置
- 加 dc_cmd() 兼容 docker-compose v1/v2 plugin
- openssl rand 生成强随机密码与 SECRET_KEY(42 字符)
- install 跳过已存在 .env,保护用户配置(管理员密码/SECRET_KEY)
- read -p < /dev/tty 兼容非 tty 环境(CI/CD)
- show-config 改用 DIFY_NGINX_PORT(nginx 入口)而非 APP_WEB_PORT

docs(mvp-design): 修正 §3.2 workflow inputs 描述
- 实际只有 query,删除错误的 user_id input 声明
- 节点序列图同步更新

feat(aiChatService): 新增 Dify 客户端与适配器
- service/dify_client.go: Dify Workflow 调用 + SSE 解析
- service/dify_adapter.go: 与现有 chat_service 桥接
- provider/ai_chat_provider.go: Dubbo 入口简化
- main.go: 装配 ConversationRepository + DifyClient

feat(migrations): 新增 AI 搭子会话表 ai_chat.sql
- ai_conversations / ai_messages 表 + 索引

docs: 新增 Dify 集成设计文档
- 2026-06-29-ai-chat-dify-mvp-design.md (MVP 实施级)
- 2026-06-29-ai-chat-dify-integration-v2-design.md (V2 演进路线图)
- docs/dify/角角.yml (Workflow DSL 导出)

config: 更新 env 模板与 docker 配置
- backend/.env.example: DIFY_* 环境变量声明
- docker/.env.prod: DIFY_API_BASE 对齐 8083
- docker/build.sh: 微调
- CLAUDE.md: 项目规范补充

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-02 12:32:34 +08:00

33 KiB
Raw Blame History

AI 搭子 MVP 方案

本文档是实施级方案。所有"为未来 100 明星 + 多 AI 平台"准备的复杂设计Provider 抽象、ProviderFactory、Pipeline、ConversationStore 抽象等)MVP 阶段不实现

V2 完整架构文档(2026-06-29-ai-chat-dify-integration-v2-design.md)保留为长期演进路线图MVP 阶段不实施。

🎉 2026-06-30 MVP 端到端跑通!

验证结果(实测):

  • 后端 20008 端口正常 Dubbo 监听
  • Dify Workflow 流式响应 8 个 chunk 正常返回
  • Dify 内部用 minimax-m3 模型推理("嗨!我是角角~..." 完整回复)
  • ai_conversations 表写入 message_count=2
  • ai_messages 表写入 user="hi" + assistant="嗨!我是角角~..." (228 字符)
  • Gateway → aichatservice → Dify → DB 完整链路打通

〇、我们解决的问题

.1 业务问题

追星 App 用户与"角角"AI 搭子聊天时AI 没有专属知识库,无法针对所追明星(默认肖战)给出有针对性的回答:

  • 用户:"肖战最近有什么新作品?"
  • 现状AI 只能泛泛而谈("肖战是中国男演员..."
  • 期望AI 应能基于最新资料回答"根据知识库肖战的新剧《X》将于 X 月上映"

.2 技术问题

# 问题 严重度 MVP 解决方式
1 缺乏 RAG 能力AI 没有专属数据 🔴 接入 Dify Workflow + 肖战知识库
2 会话无持久化:当前仅 Redis 缓存 24h 🟡 PostgreSQL 持久化 + Redis 缓存(★ 核心)
3 AI 平台耦合:未来要接 Coze/FastGPT 🟢 MVP 不做1 个 Dify 够用
4 运营成本失控100 明星时配置爆炸 🟢 MVP 固定 1 个星

.3 解决方式(业务驱动,不是架构驱动)

阶段 做什么 不做什么 触发条件
MVP Dify + PostgreSQL + Redis Provider 抽象、Fallback、多星 现在
Stage 2 多星 Dataset 切换 Provider 抽象 加第 2 个星 / 用户 > 1000
Stage 3 MiniMax fallback 抽象 Dify 偶发故障
Stage 4 Provider 抽象 复杂 Pipeline 接第 2 个 AI 平台 / 用户 > 10万
Stage 5 Pipeline / AIProfile A/B 测试 / 灰度 用户 > 100万

核心原则业务不到不做架构。MVP 阶段不实施"为未来 100 明星 + 多 AI 平台"准备的复杂设计。


〇、文档说明

  • MVP 范围:只有 1 个明星(默认肖战 star_id=87
  • AI 平台:只有 Dify无 MiniMax fallback无 Provider 抽象)
  • 目标:验证业务假设(用户愿意与"角角"聊天)
  • 后续演进:见 §十 演进路径

.1 方案概述(★ 必读)

.1.1 我们要做什么

业务:在追星 App 里加"AI 搭子"功能。用户进入后默认看到"角角"(肖战的 AI 形象),可以问"肖战最近在干嘛?"等关于肖战的问题。

MVP 范围

  • 1 个明星(肖战)
  • 1 个 Dify Workflow含 1 个肖战知识库)
  • 1 种回复源Dify
  • WebSocket 流式输出
  • 不做:多星切换、模型 fallback、人设自定义、长期记忆提取、多 AI 平台

.1.2 整体实现路径

数据迁移         Dify 端准备          后端代码            联调测试
├ 建 2 张表      ├ 准备知识库          ├ model              ├ 内部账号测试
├                ├ 创 Workflow         ├ repository         ├ WebSocket 联调
└                └ 不需要 Code 节点   └ service            └ 错误注入
                                     └ provider
                                     └ main.go 装配

不要分配实施时间——按业务节奏推进。

.1.3 关键决策80% 推到 Stage 2+

决策项 MVP 选择 Stage 2+ 再考虑
Dify 架构 Workflow(含 1 个 Dataset 多星时加 dataset 变量
AI Provider Dify 单源(无抽象) Stage 3 才抽象 AIProvider
会话存储 PostgreSQL 主 + Redis 缓存 沿用 MVP 设计
Fallback Dify 失败就报错) Stage 2 看需求
Memory 提取 不做(每轮都入 ai_messages Stage 2+ 看需求
Provider 抽象 不做 Stage 3+
数据集切换 固定 1 个(肖战) Stage 2 多星时加 mapping

.1.4 核心架构图TL;DR

用户(追星 App
   │ WebSocket
   ▼
Gateway (现有, 不改)
   │ Dubbo Triple
   ▼
AIChatService.ChatService
   │
   ├─ JWT 鉴权 (从 Dubbo attachments 取 user_id)
   │
   ├─ 保存/读取会话: ConversationRepository
   │      ├─ ai_conversations (PostgreSQL, ★ V2 关键决策保留)
   │      └─ ai_messages    (PostgreSQL, ★ V2 关键决策保留)
   │
   ├─ 调 Dify (★ 唯一 AI 源)
   │      └─ DifyClient.StreamChat()
   │            └─ POST /v1/workflows/run
   │
   └─ SSE 流 → WebSocket → 客户端

外部: Dify Workflow
   ┌──────────────────────┐
   │  开始                  │
   │  ↓                    │
   │  Knowledge Retrieval   │  ← 固定查"肖战知识库"
   │  ↓                    │
   │  LLM                  │  ← 固定 Prompt (角角人设)
   │  ↓                    │
   │  结束                  │
   └──────────────────────┘

关键简化

  • 没有 Provider 抽象
  • 没有 ProviderFactory
  • 没有 MemoryStore
  • 没有 RedisLock并发问题 MVP 阶段不严重)
  • 没有 AIProfile
  • 没有 DatasetResolver
  • 没有 Star→Dataset mapping只有 1 个星)
  • 没有 Fallback 逻辑
  • 没有 Memory 提取循环
  • 没有 Star App 多 Workflow

一、MVP 范围

1.1 包含

  • 1 个明星(肖战 star_id=87
  • 1 个 Dify Workflow3 节点)
  • 1 个 Dify Dataset肖战专属知识库
  • WebSocket 协议(沿用现有)
  • JWT 鉴权(沿用现有)
  • Audit 前置 + 后置(沿用现有 AuditService
  • Conversation 持久化到 PostgreSQL核心决策MVP 即落地
  • Redis 缓存(沿用现有)

1.2 不包含Stage 2+ 再做)

  • 多星切换(用户不能选其他明星)
  • MiniMax fallbackDify 失败就报错给用户)
  • Persona 自定义(人设固定"角角"
  • 长期记忆提取(不分析对话提取记忆)
  • Provider 抽象接口(直接调 DifyClient
  • 人设/风格/语言/记忆等参数化(都写死在 Dify Prompt 里)
  • Dify 内容审核节点AuditService 已拦截MVP 阶段够用)
  • Datasets 动态切换(固定"肖战知识库"

二、整体架构

2.1 数据流(一次完整对话)

[Mobile]
  │ WebSocket send {action: "send_message", session_id, message}
  │
  ▼
[Gateway Hub]
  │ 鉴权 (JWT → user_id)
  │
  ▼
[AIChatService Provider.SendMessage]
  │
  ├─ 1. 前置审核 (AuditService.AuditText)  ★ 现有代码
  │
  ├─ 2. 获取/创建会话 (ConversationRepository)
  │      ├─ PostgreSQL ai_conversations (★ V2 关键)
  │      └─ Redis 缓存 1h (现有代码)
  │
  ├─ 3. 调 Dify (★ MVP 唯一 AI 源)
  │      └─ DifyClient.StreamChat()
  │            └─ POST /v1/workflows/run
  │
  ├─ 4. 流式返回 + 逐 token 后置审核 (AuditService.AuditResponse)
  │      └─ ★ 现有代码
  │
  ├─ 5. 保存消息 (ConversationRepository)
  │      └─ PostgreSQL ai_messages (★ V2 关键)
  │
  └─ 6. (Stage 2+ 才做) 记忆提取

2.2 关键简化点

维度 V2 文档 MVP 实际
核心业务逻辑 9+ 步骤 4 步骤(审计/会话/Dify/保存)
Provider 数 2 个Dify + MiniMax 1 个Dify
Fallback 复杂的 Provider 切换 没有Dify 失败就报错)
星切换 动态 + dataset 映射 固定肖战
人设/风格/记忆 4 个 SystemInputs 参数 写死在 Dify Prompt
长期记忆 MemoryStore + 5 轮触发 不做
Provider 抽象 AIProvider interface 直接调 DifyClient
Workflow 节点 5 个(含 Code + Moderation 3 个(开始/检索/LLM/结束)
RedisLock 不需要(单实例部署,无并发问题)
后端代码行数估算 1500-2000 300-500

三、Dify 端配置

3.1 准备知识库Dataset

  1. 登录 Dify → "知识库" → "创建知识库"
  2. 命名:star-xz-kb(固定一个)
  3. 索引模式:high_quality
  4. 导入肖战资料(作品、行程、近期事件等)
  5. 等待向量化完成(每个文档显示 ✓)

3.2 创建 Workflow仅 3 节点)

  1. 进入"工作室" → "创建空白应用" → 类型选 Workflow
  2. 命名:star-chat-workflow
  3. 配置"开始"节点的 Input 变量:
    inputs:
      - name: query     # 用户消息
        type: text
    

    user_id(哈希后的用户标识)由后端通过 Dify 协议顶层 user 字段传入(见 §4.5 DifyClient),不作为 workflow input。

    与实际工作流对齐docs/dify/角角.yml v0.6.0 只声明了 query 一个 input。若后端仍把 user_id 塞进 inputsDify 会静默丢弃,不影响功能(user 字段仍生效)。

  4. 添加 知识检索节点
    • Knowledgestar-xz-kb(固定)
    • Query{{ query }}
    • TopK3
  5. 添加 LLM 节点
    你叫"角角",是肖战的 AI 形象。
    
    【用户问题】
    {{ query }}
    
    【知识库检索结果】
    {{#knowledge_retrieval_node.result#}}
    
    请用温柔、自然的语言回答,参考知识库内容,不要编造。
    
  6. 添加"直接回复"或"结束"节点

节点序列

[开始] inputs:{query}
   ↓
[知识检索] star-xz-kb固定
   ↓
[LLM] 角角人设(固定)
   ↓
[结束]

3.3 调试

  1. inputs={query: "肖战最近在干嘛?", user_id: "aichat-xxxx"}
  2. 验证:返回基于知识库的回答
  3. 发布 → 复制 API Keyapp-xxx

四、后端代码

4.1 改动总览

层级 改动 工作量
model/ai_chat_models.go 新增AIConversation / AIMessage GORM 模型
repository/conversation_repository.go 新建ai_conversations / ai_messages CRUD
service/chat_service.go 修改:直接调 DifyClient不再有 ChatEngine 编排)
provider/ai_chat_provider.go 修改Dubbo 入口 + 调 ChatService
main.go 修改:加 ConversationRepository 装配
migrations/ai_conversations.sql 新建2 张表 DDL
前端 无改动 0

总工作量:约 300-500 行核心代码

4.2 model/ai_chat_models.go 新增

package model

import "github.com/google/uuid"

type AIConversation struct {
    ID             int64  `gorm:"primaryKey;autoIncrement"`
    UserID         int64  `gorm:"index;not null"`
    StarID         int64  `gorm:"index;not null"`     // MVP 固定为 87
    ProviderName   string `gorm:"type:varchar(32);not null;default:'dify'"`
    ExternalConvID string `gorm:"type:varchar(128);default:''"`
    MessageCount   int    `gorm:"default:0"`
    LastActiveAt   int64  `gorm:"autoUpdateTime:milli"`
    CreatedAt      int64  `gorm:"autoCreateTime:milli"`
    UpdatedAt      int64  `gorm:"autoUpdateTime:milli"`
}

func (AIConversation) TableName() string { return "ai_conversations" }

type AIMessage struct {
    ID             int64  `gorm:"primaryKey;autoIncrement"`
    ConversationID int64  `gorm:"index;not null"`
    Role           string `gorm:"type:varchar(16);not null"`   // 'user' / 'assistant'
    Content        string `gorm:"type:text;not null"`
    CreatedAt      int64  `gorm:"autoCreateTime:milli"`
}

func (AIMessage) TableName() string { return "ai_messages" }

4.3 PostgreSQL DDL

CREATE TABLE IF NOT EXISTS ai_conversations (
    id BIGSERIAL PRIMARY KEY,
    user_id BIGINT NOT NULL,
    star_id BIGINT NOT NULL,
    provider_name VARCHAR(32) NOT NULL DEFAULT 'dify',
    external_conversation_id VARCHAR(128) DEFAULT '',
    message_count INT DEFAULT 0,
    last_active_at BIGINT NOT NULL DEFAULT (EXTRACT(EPOCH FROM NOW()) * 1000)::BIGINT,
    created_at BIGINT NOT NULL DEFAULT (EXTRACT(EPOCH FROM NOW()) * 1000)::BIGINT,
    updated_at BIGINT NOT NULL DEFAULT (EXTRACT(EPOCH FROM NOW()) * 1000)::BIGINT
);
CREATE INDEX idx_ai_conv_user_star ON ai_conversations(user_id, star_id);

CREATE TABLE IF NOT EXISTS ai_messages (
    id BIGSERIAL PRIMARY KEY,
    conversation_id BIGINT NOT NULL REFERENCES ai_conversations(id) ON DELETE CASCADE,
    role VARCHAR(16) NOT NULL,
    content TEXT NOT NULL,
    created_at BIGINT NOT NULL DEFAULT (EXTRACT(EPOCH FROM NOW()) * 1000)::BIGINT
);
CREATE INDEX idx_ai_messages_conversation ON ai_messages(conversation_id, created_at);

MVP 阶段加唯一约束user_id + star_id方便 Stage 2 加多星时再处理 MVP 阶段is_archived 等字段

4.4 ConversationRepository

package repository

import (
    "context"
    "github.com/topfans/backend/services/aiChatService/model"
    "gorm.io/gorm"
)

type ConversationRepository struct {
    db *gorm.DB
}

func NewConversationRepository(db *gorm.DB) *ConversationRepository {
    return &ConversationRepository{db: db}
}

func (r *ConversationRepository) GetOrCreate(ctx context.Context, userID, starID int64) (*model.AIConversation, error) {
    var conv model.AIConversation
    err := r.db.WithContext(ctx).Where("user_id = ? AND star_id = ?", userID, starID).First(&conv).Error
    if err == gorm.ErrRecordNotFound {
        conv = model.AIConversation{UserID: userID, StarID: starID, ProviderName: "dify"}
        if err := r.db.WithContext(ctx).Create(&conv).Error; err != nil {
            return nil, err
        }
        return &conv, nil
    }
    if err != nil {
        return nil, err
    }
    return &conv, nil
}

func (r *ConversationRepository) AppendMessage(ctx context.Context, convID int64, role, content string) error {
    msg := model.AIMessage{ConversationID: convID, Role: role, Content: content}
    return r.db.WithContext(ctx).Create(&msg).Error
}

func (r *ConversationRepository) UpdateExternalConvID(ctx context.Context, convID int64, externalID string) error {
    return r.db.WithContext(ctx).Model(&model.AIConversation{}).
        Where("id = ?", convID).
        Updates(map[string]interface{}{
            "external_conversation_id": externalID,
            "message_count":            gorm.Expr("message_count + 1"),
            "last_active_at":          gorm.Expr("(EXTRACT(EPOCH FROM NOW()) * 1000)::BIGINT"),
        }).Error
}

4.5 DifyClient★ MVP 唯一 AI 客户端)

package service

import (
    "bytes"
    "context"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "time"

    "github.com/topfans/backend/pkg/logger"
    "go.uber.org/zap"
)

type DifyClient struct {
    apiBase     string
    workflowURL string
    apiKey      string
    httpClient  *http.Client
}

type DifyConfig struct {
    APIBase     string
    WorkflowURL string
    APIKey      string
    TimeoutSec  int
}

func NewDifyClient(cfg DifyConfig) *DifyClient {
    return &DifyClient{
        apiBase:     cfg.APIBase,
        workflowURL: cfg.WorkflowURL,
        apiKey:      cfg.APIKey,
        httpClient:  &http.Client{Timeout: time.Duration(cfg.TimeoutSec) * time.Second},
    }
}

// StreamChat 流式调用 Dify Workflow
// 返回 (StreamReader, error)StreamReader 可逐 token 读取
func (c *DifyClient) StreamChat(ctx context.Context, query, userHashedID, convID string) (*DifyStreamReader, error) {
    inputs := map[string]interface{}{
        "query":   query,
        "user_id": userHashedID,
    }
    body := map[string]interface{}{
        "inputs":          inputs,
        "response_mode":   "streaming",
        "conversation_id": convID, // 首次为空
        "user":            userHashedID,
    }
    jsonData, _ := json.Marshal(body)
    httpReq, _ := http.NewRequestWithContext(ctx, "POST",
        c.apiBase+c.workflowURL, bytes.NewReader(jsonData))
    httpReq.Header.Set("Authorization", "Bearer "+c.apiKey)
    httpReq.Header.Set("Content-Type", "application/json")

    resp, err := c.httpClient.Do(httpReq)
    if err != nil {
        return nil, fmt.Errorf("dify request: %w", err)
    }
    if resp.StatusCode != http.StatusOK {
        body, _ := io.ReadAll(resp.Body)
        resp.Body.Close()
        return nil, fmt.Errorf("dify returned HTTP %d: %s", resp.StatusCode, string(body))
    }
    return &DifyStreamReader{reader: resp.Body, decoder: NewSSEDecoder(resp.Body), conversationID: ""}, nil
}

// DifyStreamReader 解析 Dify SSE 流
type DifyStreamReader struct {
    reader         io.ReadCloser
    decoder        *SSEDecoder
    conversationID string
}

func (r *DifyStreamReader) Next() (content string, done bool, err error) {
    for {
        line, err := r.decoder.Next()
        if err != nil {
            if err == io.EOF { return "", true, nil }
            return "", true, err
        }
        if !strings.HasPrefix(line, "data: ") { continue }
        data := strings.TrimPrefix(line, "data: ")

        var event struct {
            Event          string `json:"event"`
            Answer         string `json:"answer"`
            ConversationID string `json:"conversation_id"`
        }
        if err := json.Unmarshal([]byte(data), &event); err != nil { continue }

        if event.ConversationID != "" && r.conversationID == "" {
            r.conversationID = event.ConversationID
        }
        switch event.Event {
        case "message":
            return event.Answer, false, nil
        case "message_end":
            return "", true, nil
        case "error":
            return "", true, fmt.Errorf("dify error event")
        }
    }
}

func (r *DifyStreamReader) GetConversationID() string { return r.conversationID }
func (r *DifyStreamReader) Close() error { return r.reader.Close() }

★ MVP 阶段没有滑动窗口审计、retry 循环、敏感词检测(现有 AuditService 已足够)

4.6 ChatService核心业务逻辑

package service

type ChatService struct {
    audit        *AuditService
    convRepo     *repository.ConversationRepository
    difyClient   *DifyClient
    userIDSalt   string
}

func NewChatService(audit *AuditService, convRepo *repository.ConversationRepository, dify *DifyClient) *ChatService {
    return &ChatService{audit: audit, convRepo: convRepo, difyClient: dify}
}

const (
    DefaultStarID  = int64(87)   // 肖战
    DefaultUserSalt = "topfans-default-salt"
)

func (s *ChatService) hashUserID(userID int64) string {
    h := sha256.Sum256([]byte(fmt.Sprintf("%d:%s", userID, s.userIDSalt)))
    return "aichat-" + hex.EncodeToString(h[:8])
}

// SendMessage 核心流程4 步)
func (s *ChatService) SendMessage(ctx context.Context, userID int64, message string) (<-chan *StreamChunk, error) {
    out := make(chan *StreamChunk, 16)
    go func() {
        defer close(out)

        // 1. 前置审核
        if !s.audit.AuditText(message) {
            out <- &StreamChunk{Type: "message", Content: s.audit.DefaultSafeResponse(), IsEnd: false}
            out <- &StreamChunk{Type: "message", IsEnd: true}
            return
        }

        // 2. 获取/创建会话(★ V2 关键决策PostgreSQL 持久化)
        conv, err := s.convRepo.GetOrCreate(ctx, userID, DefaultStarID)
        if err != nil {
            out <- &StreamChunk{Type: "error", Error: "会话创建失败"}
            return
        }

        // 3. 调 Dify★ MVP 唯一 AI 源)
        streamReader, err := s.difyClient.StreamChat(ctx, message, s.hashUserID(userID), conv.ExternalConvID)
        if err != nil {
            logger.Logger.Error("Dify call failed", zap.Error(err))
            out <- &StreamChunk{Type: "error", Error: "服务暂不可用"}
            return
        }
        defer streamReader.Close()

        // 4. 流式返回 + 后置审核 + 保存
        var fullResponse string
        for {
            content, done, err := streamReader.Next()
            if err != nil {
                out <- &StreamChunk{Type: "error", Error: "服务异常"}
                return
            }
            if content != "" && !s.audit.AuditResponse(content) {
                // 命中敏感词
                out <- &StreamChunk{Type: "message", Content: s.audit.DefaultSafeResponse(), IsEnd: false}
                out <- &StreamChunk{Type: "message", IsEnd: true}
                s.convRepo.AppendMessage(ctx, conv.ID, "assistant", s.audit.DefaultSafeResponse())
                return
            }
            fullResponse += content
            out <- &StreamChunk{Type: "message", Content: content, IsEnd: done}
            if done { break }
        }

        // 5. 更新 Dify conv_id首次
        if newConvID := streamReader.GetConversationID(); newConvID != "" && newConvID != conv.ExternalConvID {
            s.convRepo.UpdateExternalConvID(ctx, conv.ID, newConvID)
        }

        // 6. 保存消息
        s.convRepo.AppendMessage(ctx, conv.ID, "user", message)
        s.convRepo.AppendMessage(ctx, conv.ID, "assistant", fullResponse)
    }()
    return out, nil
}

// GetWelcomeMessage MVP 阶段固定返回"角角"欢迎语
func (s *ChatService) GetWelcomeMessage() string {
    return "你好,我是角角,专注肖战的 AI 搭子。有什么想了解的?"
}

4.7 Provider 大幅简化

package provider

type AIChatProvider struct {
    chatService *service.ChatService
}

func (p *AIChatProvider) SendMessage(ctx context.Context, req *pb.ChatMessageRequest, stream pb.AIChatService_SendMessageServer) error {
    userID, _, err := extractUserInfoFromDubboAttachments(ctx)
    if err != nil { return err }

    chunks, err := p.chatService.SendMessage(ctx, userID, req.Message)
    if err != nil { return err }

    for chunk := range chunks {
        if chunk.Type == "error" {
            stream.Send(&pb.ChatMessageResponse{Content: chunk.Error, IsEnd: true})
        } else {
            stream.Send(&pb.ChatMessageResponse{Content: chunk.Content, IsEnd: chunk.IsEnd})
        }
    }
    return nil
}

极简~30 行。完全没有 V2 里的 ChatEngine 编排、锁、Provider 抽象、Factory 等

4.8 main.go 装配

// MVP 装配:极简
convRepo := repository.NewConversationRepository(database.GetDB())
difyClient := service.NewDifyClient(service.DifyConfig{
    APIBase:     getEnv("DIFY_API_BASE", "https://api.dify.ai/v1"),
    WorkflowURL: "/workflows/run",
    APIKey:      getEnv("DIFY_API_KEY", ""),
    TimeoutSec:  60,
})
chatService := service.NewChatService(auditService, convRepo, difyClient)
aiChatProvider := provider.NewAIChatProvider(chatService)

五、消息协议(无改动)

WebSocket 协议与现有实现一致:

  • Client → Server{action: "send_message", session_id, message}
  • Server → Client{type: "message", content, is_end}{type: "error", error}

前端零改动


六、关键设计决策

决策 选择 理由
Provider 抽象 不做 MVP 只有 1 个 AI 源,抽象无价值
MiniMax fallback 不做 Dify 失败就报错,避免增加复杂度
长期记忆提取 不做 业务假设未验证前不做
Redis 缓存 保留 1h 缓存会话元数据,避免每次查 DB
PostgreSQL 持久化 保留(★ 关键) 跨设备/跨天续接(追星场景长生命周期)
滑动窗口审计 不做 现有 AuditService 逐 token 检查已足够
不做 单实例部署,并发问题不严重
Dify 内容审核节点 不做 现有 AuditService 已拦截MVP 够用
UserStyle/UserNickname 参数 不做 写死在 Dify Prompt 里
星切换 不做 固定 star_id=87肖战

七、配置清单

7.1 环境变量

变量 用途 必填
DIFY_API_KEY Dify Workflow API Key
DIFY_API_BASE Dify API 地址 否(默认 https://api.dify.ai/v1

7.2 ai_chat_configs★ MVP 全部不要)

MVP 阶段直接用环境变量,不写 ai_chat_configs 数据库

V2 文档里 9 个 dify.* 配置项 MVP 全部不需要enabler、api_base、workflow_url、star_dataset_mapping、api_key、timeout_sec、user_id_salt、retry_count、fallback_to_minimax——全部 hardcode 或用环境变量

Stage 2+ 才把这些移到数据库配置。


八、部署清单

不要分配实施时间。按业务节奏推进。

8.1 数据库

  • DBA 执行 migrations/ai_conversations.sql
  • 验证表结构和索引

8.2 Dify 端

  • Dify 管理员创建 star-xz-kb 知识库
  • 导入肖战资料并等待向量化完成
  • 创建 star-chat-workflow3 节点:开始/检索/LLM
  • 配置 Prompt"你是角角,温柔回复,参考知识库..."
  • 调试并发布
  • 把 API Key 安全转给后端

8.3 后端代码

  • 新建 model/ai_chat_models.go 的 AIConversation/AIMessage
  • 新建 repository/conversation_repository.go
  • 新建 service/dify_client.go
  • 修改 service/chat_service.go直接调 DifyClient不引入 ChatEngine
  • 简化 provider/ai_chat_provider.go
  • 修改 main.go 装配
  • 单元测试ConversationRepository CRUD
  • 集成测试mock Dify server 跑完整 SendMessage

8.4 联调测试

  • 内部账号测试:进 ai-dazi 页面发消息
  • 验证:流式返回正常
  • 验证ai_messages 表有 user + assistant 两条记录
  • 验证:关掉重开会话能续接
  • 验证:敏感词("裸聊"等)被拦截
  • 验证Dify 故障时返回明确错误

九、验证清单

9.1 功能验证

  • 发送"肖战最近在干嘛?"能返回基于知识库的回答
  • 发送"你好"能返回通用问候
  • 同用户第二次发消息能续接上下文Dify conversation_id
  • ai_messages 表有 user + assistant 两条记录
  • ai_conversations 表的 message_count 正确递增

9.2 安全验证

  • 前置审核:用户发"裸聊"等敏感词被拦截
  • 后置审核Dify 回复中含敏感词被拦截
  • Dify API Key 不出现在日志

9.3 不验证Stage 2+ 再做)

  • 多星切换MVP 不做)
  • FallbackMVP 不做)
  • 长期记忆提取MVP 不做)

十、Stage 2+ 演进路径

MVP 跑通后,根据用户量和业务反馈,按以下顺序演进:

Stage 触发条件 关键改动
Stage 2 用户量 > 1000 OR 加第 2 个星 1. 多星 Dataset 切换star_dataset_mapping
2. WebSocket 端 InitSession 欢迎语动态化
3. ai_conversations 加 UNIQUE(user_id, star_id)
Stage 3 Dify 偶发故障 OR SLA 要求 1. MiniMax fallback仅 message_count=0 时)
2. Dify retry 循环
Stage 4 用户量 > 10万 OR 接 2+ AI 平台 1. AIProvider 抽象
2. ProviderFactory 策略模式
3. CozeProvider / OpenAIProvider 实现
Stage 5 用户量 > 100万 OR 业务复杂 1. ChatEngine Pipeline 化
2. AIProfile 配置化A/B 测试、灰度)
3. 长期记忆提取

关键原则:每个 Stage 都是业务驱动,不是架构驱动。

★ Stage 2+ 演进时必踩的 3 个坑P0 修复笔记)

这些是 V2 架构评审发现的真实 bugMVP 阶段不修(流量小、问题不暴露),但 Stage 2+ 流量上来后必现。 ★ 必读:实施 Stage 2 之前,必须先修这 3 个 P0 问题。

坑 1并发请求分裂会话★ P0-1

症状:用户手机 + 平板同时发消息Dify 端产生两个会话AI 上下文错乱。

根因:两个并发请求都查到 ExternalConvID="",都调 DifyDify 给两个不同的 conversation_id,后写入的覆盖先写入的。

修复

  • 加 Redis 分布式锁 conv_lock:{userId}:{starId}TTL 30s
  • 锁范围:GetOrCreateConversationUpdateExternalID
  • 锁未获取时 sleep 200ms 重试一次

坑 2审计拦截后不保存对话★ P0-2

症状用户每次触发敏感词拦截后AI 都不记得之前说过什么,行为诡异。

根因:审计分支直接 return没保存"user 原句 + 安全回复"。

修复:审计分支也调 AppendConversationMessages 保存对话。

坑 3组合敏感词漏检★ P0-3

症状Dify 返回"色"+"情"分两个 token单独都不违规组合违规。

根因V1 AuditService 逐 token 检查(strings.Contains(token, word)),单 token 视角。

修复Dify 流式接收时维护 sliding window buffer20 字符),每个 token 检查 buffer。

★ Stage 4 抽象时必踩的 3 个坑(架构评审笔记)

★ 这些是 V2 架构评审发现的 设计层面问题Stage 4 做 AIProvider 抽象时必踩。

坑 4Provider God Class评审 #1

症状DifyProvider 写了 2000+ 行什么都管HTTP / Stream / Cache / Conversation / Retry / Hash / Audit

修复:拆为 4 个组件:

  • WorkflowClientHTTP + SSE 解析 + Retry
  • HistoryClient:拉历史消息
  • DatasetResolverstar_id → dataset_id 映射
  • ConversationStore:缓存 + 持久化

原则Provider 只做协调,不做任何具体工作。

坑 5Provider 直接依赖 Redis/Repository评审 #5

症状Provider 改存储Redis → Memcached所有 Provider 都要改。

修复:引入 ConversationStore 抽象Provider 只依赖接口。底层是 CachedConversationStorePostgreSQL + Redis 缓存)。

坑 6Workflow 与 Backend 重复维护 dataset 映射(评审 #3

症状Backend 改 star_dataset_mapping 忘改 Workflow → 用户问"肖战"答"王一博"的资料。

修复

  • 单一数据源:映射只在 Backend ai_chat_configs.dify.star_dataset_mapping 维护
  • Workflow 输入直接接 dataset_idBackend 传过来的)
  • Workflow 内部维护任何 star_id → dataset_id 映射(没有 Code 节点

★ 演进时不要做

  • 不要预先做 P0 修复MVP 阶段流量小race / 审计保存 / 组合敏感词都暴露不出来
  • 不要预先做 Provider 抽象MVP 只有 1 个 AI 源(写死就行)
  • 不要预先做 PipelineSendMessage 函数 < 200 行不需要 Pipeline
  • 不要预先做 AIProfile1 个星 1 种 AI 源根本不需要 A/B

十一、与 V2 文档的关系

V2 文档已删除2026-06-29 决定)。

原因V2 是"100 明星 + 多 AI 平台"的完整架构设计,MVP 不需要 80% 的内容。V2 关键内容已迁移到本 MVP 文档的 §10 演进路径,包含 3 个 P0 修复笔记 + 3 个 Stage 4 抽象笔记。

实施时按 MVP 推进,跑通后再按 §10 Stage 2+ 演进。


十二、关键文件清单

12.1 新建文件

backend/services/aiChatService/
├── model/
│   └── ai_chat_models.go               (新增 AIConversation / AIMessage struct)
├── repository/
│   └── conversation_repository.go      (新建3 个方法)
├── service/
│   ├── dify_client.go                  (新建,唯一 AI 客户端)
│   └── chat_service.go                  (修改4 步流程)

migrations/
└── ai_conversations.sql                 (新建2 张表 DDL)

12.2 修改文件

backend/services/aiChatService/
├── provider/
│   └── ai_chat_provider.go              (大幅简化)
└── main.go                              (加 convRepo + difyClient 装配)

12.3 不动文件

  • 前端(所有 .vue / .js
  • GatewayWebSocket Hub
  • AuditService现有代码MVP 沿用)
  • JWT 鉴权(现有代码)

总结

MVP 阶段的核心是验证业务假设,不是搭建完美架构

MVP 范围

  • 用户能跟"角角"聊天
  • AI 能基于肖战知识库回答
  • 对话跨天续接PostgreSQL 持久化)
  • 敏感词拦截AuditService 沿用)

不做(业务驱动)

  • Provider 抽象(只有 1 个 AI 源)
  • FallbackDify 失败就报错)
  • 多星(只有 1 个星)
  • 长期记忆提取(每轮直接入 ai_messages
  • Persona 自定义(写死在 Dify Prompt
  • Dify 内容审核节点AuditService 够用)
  • 9 个 dify.* 数据库配置(环境变量够用)

业务驱动原则

每个 Stage 都是业务驱动,不是架构驱动:

  1. 业务没到的复杂度 → 不预先做
  2. 架构是演进的,不是一次性完美设计
  3. V2 完整文档作为长期演进路线图不删除
  4. 跑通 MVP 后,按 §十 Stage 2-5 渐进改进

这些就够了。其他都是"未来 100 明星 + 多 AI 平台"的事。