package service import ( "bufio" "bytes" "context" "encoding/json" "fmt" "io" "net/http" "strings" "time" "github.com/topfans/backend/pkg/logger" "go.uber.org/zap" ) // DifyConfig Dify 客户端配置 // MVP 阶段:直接从环境变量读 // Stage 2+:移到 ai_chat_configs 数据库 type DifyConfig struct { APIBase string APIKey string TimeoutSec int } // DifyClient Dify HTTP 客户端(MVP 唯一 AI 源) // ★ 简化:只调 Workflow /v1/workflows/run 流式接口 // 不做:重试循环(V1 Network 错误重试 1 次,MVP 失败直接报错) // 不做:滑动窗口审计(V1 AuditService 逐 token 拦截已足够) type DifyClient struct { apiBase string apiKey string httpClient *http.Client } // NewDifyClient 创建 Dify 客户端 func NewDifyClient(cfg DifyConfig) *DifyClient { return &DifyClient{ apiBase: cfg.APIBase, apiKey: cfg.APIKey, httpClient: &http.Client{Timeout: time.Duration(cfg.TimeoutSec) * time.Second}, } } // StreamRequest 流式调用入参 type StreamRequest struct { Query string ConversationID string } // StreamChat 流式调用 Dify Workflow // 输入:用户消息 + Dify conversation_id(首次为空) // 输出:返回 (DifyStreamReader, error) func (c *DifyClient) StreamChat(ctx context.Context, req StreamRequest) (*DifyStreamReader, error) { body := map[string]interface{}{ "inputs": map[string]string{"query": req.Query}, "response_mode": "streaming", "conversation_id": req.ConversationID, "user": "aichat-mvp", } jsonData, err := json.Marshal(body) if err != nil { return nil, fmt.Errorf("dify request marshal: %w", err) } url := c.apiBase + "/workflows/run" httpReq, err := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(jsonData)) if err != nil { return nil, fmt.Errorf("dify request new: %w", err) } 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 { respBody, _ := io.ReadAll(resp.Body) resp.Body.Close() return nil, fmt.Errorf("dify returned HTTP %d: %s", resp.StatusCode, string(respBody)) } logger.Logger.Info("Dify stream started", zap.String("url", url)) return newDifyStreamReader(resp.Body), nil } // ============================================================ // DifyStreamReader:解析 SSE 流 // ============================================================ // newDifyStreamReader 创建流式读取器 func newDifyStreamReader(body io.ReadCloser) *DifyStreamReader { scanner := bufio.NewScanner(body) // 增大 buffer 防止 Dify 流式 chunk 过大被截断 scanner.Buffer(make([]byte, 64*1024), 1024*1024) return &DifyStreamReader{ body: body, conversationID: "", closed: false, scanner: scanner, } } // DifyStreamReader Dify SSE 流式读取器 // 简化版本:用 bufio.Scanner 按行读 data: 字段 type DifyStreamReader struct { body io.ReadCloser conversationID string closed bool scanner *bufio.Scanner } // Next 读下一条数据 // 返回 (content, done, convID): // - content: 当前 token(流式) // - done: 流是否结束 // - convID: 流结束后从 convID 读 Dify conversation_id func (r *DifyStreamReader) Next() (content string, done bool, convID string) { if r.closed { return "", true, r.conversationID } for r.scanner.Scan() { line := strings.TrimSpace(r.scanner.Text()) if line == "" { continue } if !strings.HasPrefix(line, "data: ") { continue } payload := strings.TrimPrefix(line, "data: ") if payload == "[DONE]" { r.closed = true return "", true, r.conversationID } // ★ MVP 修复:Dify Workflow 端点 (/v1/workflows/run) 的事件格式 // - text_chunk: 流式分块,text 字段是内容 // - workflow_finished: 结束事件,data.outputs.text 是完整内容 // - 没有 conversation_id(Workflow 端点不维护 conversation_id) // ★ 与 Chatflow 端点 (/v1/chat-messages) 不同: // - Chatflow 用 event=message, answer 字段 // - Workflow 用 event=text_chunk, text 字段 var event struct { Event string `json:"event"` Text string `json:"text"` // Chatflow: 顶层 text ConversationID string `json:"conversation_id"` Message string `json:"message"` Data struct { Text string `json:"text"` // Workflow: data.text } `json:"data"` } if err := json.Unmarshal([]byte(payload), &event); err != nil { continue } // 记录 conversation_id(仅 Chatflow 有;Workflow 端点通常无) if event.ConversationID != "" && r.conversationID == "" { r.conversationID = event.ConversationID } text := event.Data.Text if text == "" { text = event.Text } switch event.Event { case "text_chunk": if text != "" { return text, false, "" } case "message": if text != "" { return text, false, "" } case "workflow_finished": // ★ MVP 关键:Dify Workflow 结束事件 r.closed = true return "", true, r.conversationID case "message_end": // Chatflow 端点兼容 r.closed = true return "", true, r.conversationID case "error": logger.Logger.Error("Dify stream error event", zap.String("message", event.Message)) r.closed = true return "", true, r.conversationID } } // scanner 结束(流关闭或出错) r.closed = true if err := r.scanner.Err(); err != nil { logger.Logger.Warn("Dify stream scan error", zap.Error(err)) } return "", true, r.conversationID } // Close 关闭流 func (r *DifyStreamReader) Close() error { if r.closed { return nil } r.closed = true return r.body.Close() } // GetConversationID 获取 Dify 返回的 conversation_id // 仅在 Next() 返回 done=true 后调用 func (r *DifyStreamReader) GetConversationID() string { return r.conversationID }