feat(captain): align task payload persistence

This commit is contained in:
2026-06-05 13:38:55 +08:00
parent 3263ed9284
commit 13cb750a2a
11 changed files with 822 additions and 89 deletions
@@ -30,6 +30,22 @@ func (h *CaptainTaskExtendedHandler) LabelSuggestion(c *gin.Context) {
return
}
if c.Request.Method == http.MethodPost {
var req service.ChatwootLabelSuggestionRequest
if err := c.ShouldBindJSON(&req); err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid request body: "+err.Error())
return
}
result, err := h.svc.LabelSuggestion(c.Request.Context(), uint(accountID), &req)
if err != nil {
applogger.L().Errorf("Label suggestion: %v", err)
renderCaptainTaskError(c, err)
return
}
renderCaptainExtendedTaskResult(c, result)
return
}
// Parse conversation_ids from query param (comma-separated)
convIDsStr := c.Query("conversation_ids")
if convIDsStr == "" {
@@ -69,6 +85,22 @@ func (h *CaptainTaskExtendedHandler) FollowUp(c *gin.Context) {
return
}
if c.Request.Method == http.MethodPost {
var req service.ChatwootFollowUpRequest
if err := c.ShouldBindJSON(&req); err != nil {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid request body: "+err.Error())
return
}
result, err := h.svc.FollowUp(c.Request.Context(), uint(accountID), &req)
if err != nil {
applogger.L().Errorf("Follow-up: %v", err)
renderCaptainTaskError(c, err)
return
}
renderCaptainExtendedTaskResult(c, result)
return
}
convIDsStr := c.Query("conversation_ids")
if convIDsStr == "" {
response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, "conversation_ids required")
@@ -98,6 +130,18 @@ func (h *CaptainTaskExtendedHandler) FollowUp(c *gin.Context) {
response.OK(c, result)
}
func renderCaptainExtendedTaskResult(c *gin.Context, result *service.ChatwootTaskResult) {
if result == nil || result.Message == nil {
c.JSON(http.StatusOK, gin.H{"message": nil})
return
}
payload := gin.H{"message": *result.Message}
if result.FollowUpContext != nil {
payload["follow_up_context"] = result.FollowUpContext
}
c.JSON(http.StatusOK, payload)
}
// parseUintSlice parses a comma-separated string of uint values.
func parseUintSlice(s string) ([]uint, error) {
parts := strings.Split(s, ",")
@@ -114,4 +158,4 @@ func parseUintSlice(s string) ([]uint, error) {
result = append(result, uint(v))
}
return result, nil
}
}
@@ -1,6 +1,7 @@
package v1
import (
"bytes"
"context"
"encoding/json"
"fmt"
@@ -47,12 +48,14 @@ func setupTaskExtendedHandlerTest(t *testing.T, mockLLM llm.Provider) (*CaptainT
&model.Message{},
&model.Inbox{},
&model.Contact{},
&model.CopilotSuggestionMessage{},
))
convRepo := repository.NewConversationRepo(db)
msgRepo := repository.NewMessageRepo(db)
assistantRepo := repository.NewCaptainAssistantRepo(db)
prefRepo := repository.NewCaptainPreferenceRepo(db)
svc := service.NewCaptainTaskExtendedService(convRepo, msgRepo, assistantRepo, prefRepo, mockLLM)
suggestionRepo := repository.NewCopilotSuggestionRepo(db)
svc := service.NewCaptainTaskExtendedService(convRepo, msgRepo, assistantRepo, prefRepo, mockLLM, suggestionRepo)
handler := NewCaptainTaskExtendedHandler(svc)
return handler, db
}
@@ -103,6 +106,40 @@ func TestCaptainTaskExtendedHandler_LabelSuggestion_MissingConversationIDs(t *te
assert.Equal(t, http.StatusBadRequest, w.Code)
}
func TestCaptainTaskExtendedHandler_LabelSuggestion_ChatwootPostRawPayload(t *testing.T) {
mockLLM := &mockTaskExtHandlerLLM{
response: &llm.ChatResponse{Choices: []llm.ChatChoice{{Message: llm.ChatMessage{Role: "assistant", Content: "billing, urgent"}}}},
}
handler, db := setupTaskExtendedHandlerTest(t, mockLLM)
displayID := uint(77)
inbox := &model.Inbox{AccountID: 1, Name: "Inbox", ChannelType: "web_widget"}
require.NoError(t, db.Create(inbox).Error)
conv := &model.Conversation{AccountID: 1, InboxID: inbox.ID, DisplayID: &displayID, Status: "open", ChannelType: "web_widget", Channel: "web_widget"}
require.NoError(t, db.Create(conv).Error)
msg := &model.Message{ConversationID: conv.ID, AccountID: 1, InboxID: inbox.ID, SenderType: "contact", Content: "billing help", ContentType: "text", MessageType: "incoming"}
require.NoError(t, db.Create(msg).Error)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Params = gin.Params{{Key: "id", Value: "1"}}
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/accounts/1/captain/tasks/label_suggestion", bytes.NewReader([]byte(`{"conversation_display_id":77}`)))
c.Request.Header.Set("Content-Type", "application/json")
handler.LabelSuggestion(c)
require.Equal(t, http.StatusOK, w.Code)
var resp map[string]interface{}
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp))
assert.Equal(t, "billing, urgent", resp["message"])
assert.NotContains(t, resp, "success")
require.Contains(t, resp, "follow_up_context")
var stored []model.CopilotSuggestionMessage
require.NoError(t, db.Find(&stored).Error)
require.Len(t, stored, 1)
assert.Equal(t, "billing, urgent", stored[0].Content)
}
func TestCaptainTaskExtendedHandler_FollowUp(t *testing.T) {
mockLLM := &mockTaskExtHandlerLLM{
response: &llm.ChatResponse{
@@ -143,3 +180,42 @@ func TestCaptainTaskExtendedHandler_FollowUp_MissingConversationIDs(t *testing.T
assert.Equal(t, http.StatusBadRequest, w.Code)
}
func TestCaptainTaskExtendedHandler_FollowUp_ChatwootPostUpdatesContext(t *testing.T) {
mockLLM := &mockTaskExtHandlerLLM{
response: &llm.ChatResponse{Choices: []llm.ChatChoice{{Message: llm.ChatMessage{Role: "assistant", Content: "Refined answer"}}}},
}
handler, db := setupTaskExtendedHandlerTest(t, mockLLM)
displayID := uint(88)
inbox := &model.Inbox{AccountID: 1, Name: "Inbox", ChannelType: "web_widget"}
require.NoError(t, db.Create(inbox).Error)
conv := &model.Conversation{AccountID: 1, InboxID: inbox.ID, DisplayID: &displayID, Status: "open", ChannelType: "web_widget", Channel: "web_widget"}
require.NoError(t, db.Create(conv).Error)
body := []byte(`{
"conversation_display_id": 88,
"message": "Make it warmer",
"follow_up_context": {
"event_name": "professional",
"original_context": "Original draft",
"last_response": "Previous answer",
"conversation_history": []
}
}`)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Params = gin.Params{{Key: "id", Value: "1"}}
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/accounts/1/captain/tasks/follow_up", bytes.NewReader(body))
c.Request.Header.Set("Content-Type", "application/json")
handler.FollowUp(c)
require.Equal(t, http.StatusOK, w.Code)
var resp map[string]interface{}
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp))
assert.Equal(t, "Refined answer", resp["message"])
ctx := resp["follow_up_context"].(map[string]interface{})
assert.Equal(t, "Refined answer", ctx["last_response"])
history := ctx["conversation_history"].([]interface{})
require.Len(t, history, 2)
}
@@ -45,11 +45,11 @@ func (h *CaptainTaskHandler) ReplySuggestion(c *gin.Context) {
result, err := h.svc.ReplySuggestion(c.Request.Context(), uint(accountID), &req)
if err != nil {
applogger.L().Errorf("ReplySuggestion: %v", err)
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to generate reply suggestions")
renderCaptainTaskError(c, err)
return
}
response.OK(c, result)
renderCaptainTaskPayload(c, result.Message, result.FollowUpContext)
}
// Summarize generates a concise summary of a conversation.
@@ -70,11 +70,11 @@ func (h *CaptainTaskHandler) Summarize(c *gin.Context) {
result, err := h.svc.Summarize(c.Request.Context(), uint(accountID), &req)
if err != nil {
applogger.L().Errorf("Summarize: %v", err)
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to summarize conversation")
renderCaptainTaskError(c, err)
return
}
response.OK(c, result)
renderCaptainTaskPayload(c, result.Message, result.FollowUpContext)
}
// Rewrite rewrites a draft message to improve tone, clarity, or language.
@@ -95,11 +95,11 @@ func (h *CaptainTaskHandler) Rewrite(c *gin.Context) {
result, err := h.svc.Rewrite(c.Request.Context(), uint(accountID), &req)
if err != nil {
applogger.L().Errorf("Rewrite: %v", err)
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to rewrite message")
renderCaptainTaskError(c, err)
return
}
response.OK(c, result)
renderCaptainTaskPayload(c, result.Message, result.FollowUpContext)
}
// --- SSE Streaming Endpoints (M12) ---
@@ -280,3 +280,20 @@ func captainEscapeJSONString(s string) string {
}
return result.String()
}
func renderCaptainTaskPayload(c *gin.Context, message string, followUpContext map[string]interface{}) {
payload := gin.H{"message": message}
if followUpContext != nil {
payload["follow_up_context"] = followUpContext
}
c.JSON(http.StatusOK, payload)
}
func renderCaptainTaskError(c *gin.Context, err error) {
status, message, ok := service.CaptainTaskErrorStatus(err)
if !ok {
status = http.StatusUnprocessableEntity
message = err.Error()
}
c.JSON(status, gin.H{"error": message})
}
@@ -0,0 +1,103 @@
package v1
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"github.com/gin-gonic/gin"
"github.com/gochat/gochat/internal/llm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/gochat/gochat/internal/service"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
type mockCaptainTaskHandlerLLM struct {
response *llm.ChatResponse
err error
}
func (m *mockCaptainTaskHandlerLLM) ChatCompletion(_ context.Context, _ llm.ChatRequest) (*llm.ChatResponse, error) {
return m.response, m.err
}
func (m *mockCaptainTaskHandlerLLM) CreateEmbedding(_ context.Context, _ llm.EmbeddingRequest) (*llm.EmbeddingResponse, error) {
return nil, nil
}
func (m *mockCaptainTaskHandlerLLM) ChatCompletionStream(_ context.Context, _ llm.ChatRequest, _ func(llm.StreamChunk) error) error {
return nil
}
func setupCaptainTaskHandlerTest(t *testing.T, provider llm.Provider) (*CaptainTaskHandler, *gorm.DB) {
t.Helper()
gin.SetMode(gin.TestMode)
dbName := fmt.Sprintf("file:%s?mode=memory&cache=private", t.Name())
db, err := gorm.Open(sqlite.Open(dbName), &gorm.Config{})
require.NoError(t, err)
require.NoError(t, db.AutoMigrate(
&model.Account{},
&model.Conversation{},
&model.Message{},
&model.CaptainAssistant{},
&model.CaptainAssistantResponse{},
&model.CaptainCustomTool{},
&model.CaptainDocument{},
&model.CopilotSuggestionMessage{},
))
assistantRepo := repository.NewCaptainAssistantRepo(db)
responseRepo := repository.NewCaptainAssistantResponseRepo(db)
customToolRepo := repository.NewCaptainCustomToolRepo(db)
conversationRepo := repository.NewConversationRepo(db)
messageRepo := repository.NewMessageRepo(db)
suggestionRepo := repository.NewCopilotSuggestionRepo(db)
svc := service.NewCaptainTaskService(assistantRepo, responseRepo, customToolRepo, conversationRepo, messageRepo, provider, nil, suggestionRepo)
return NewCaptainTaskHandler(svc), db
}
func TestCaptainTaskHandler_Summarize_ChatwootRawPayload(t *testing.T) {
provider := &mockCaptainTaskHandlerLLM{response: &llm.ChatResponse{Choices: []llm.ChatChoice{{Message: llm.ChatMessage{Role: "assistant", Content: "Short summary"}}}}}
handler, db := setupCaptainTaskHandlerTest(t, provider)
displayID := uint(123)
conv := &model.Conversation{AccountID: 1, DisplayID: &displayID, Status: "open", ChannelType: "web_widget", Channel: "web_widget"}
require.NoError(t, db.Create(conv).Error)
require.NoError(t, db.Create(&model.Message{ConversationID: conv.ID, AccountID: 1, SenderType: "contact", MessageType: "incoming", Content: "Need help"}).Error)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Params = gin.Params{{Key: "id", Value: "1"}}
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/accounts/1/captain/tasks/summarize", bytes.NewReader([]byte(`{"conversation_display_id":123}`)))
c.Request.Header.Set("Content-Type", "application/json")
handler.Summarize(c)
require.Equal(t, http.StatusOK, w.Code)
var resp map[string]interface{}
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp))
assert.Equal(t, "Short summary", resp["message"])
assert.NotContains(t, resp, "success")
}
func TestCaptainTaskHandler_Rewrite_NoProviderRawDisabled(t *testing.T) {
handler, _ := setupCaptainTaskHandlerTest(t, nil)
w := httptest.NewRecorder()
c, _ := gin.CreateTestContext(w)
c.Params = gin.Params{{Key: "id", Value: "1"}}
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/accounts/1/captain/tasks/rewrite", bytes.NewReader([]byte(`{"content":"hello","operation":"professional"}`)))
c.Request.Header.Set("Content-Type", "application/json")
handler.Rewrite(c)
require.Equal(t, http.StatusUnprocessableEntity, w.Code)
var resp map[string]interface{}
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp))
assert.Equal(t, "Captain is disabled", resp["error"])
assert.NotContains(t, resp, "success")
}