* feat(conversations): complete manual AI takeover * fix(conversations): align AI takeover flow with channel AI * fix(conversations): close takeover review gaps --------- Co-authored-by: Rogee <rogee@ipao.vip>
219 lines
12 KiB
Go
219 lines
12 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/gochat/gochat/internal/channel"
|
|
"github.com/gochat/gochat/internal/model"
|
|
"github.com/gochat/gochat/internal/repository"
|
|
"github.com/gochat/gochat/internal/worker"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"gorm.io/driver/sqlite"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
func setupCaptainConversationWorkerTest(t *testing.T) (*gorm.DB, *CaptainConversationService, *MessageService, *model.Account, *model.Inbox, *model.Conversation, *model.CaptainAssistant) {
|
|
t.Helper()
|
|
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.Inbox{}, &model.Contact{}, &model.Conversation{}, &model.Message{}, &model.Attachment{}, &model.AgentBot{}, &model.AgentBotInbox{}, &model.CaptainAssistant{}, &model.CaptainInbox{}, &model.CaptainPreference{}, &model.BackgroundJob{}))
|
|
t.Cleanup(func() {
|
|
sqlDB, _ := db.DB()
|
|
sqlDB.Close()
|
|
})
|
|
account := &model.Account{Name: "Captain Org", Active: true}
|
|
require.NoError(t, db.Create(account).Error)
|
|
inbox := &model.Inbox{AccountID: account.ID, Name: "API", ChannelType: "api", ChannelID: 1}
|
|
require.NoError(t, db.Create(inbox).Error)
|
|
contact := &model.Contact{AccountID: account.ID, Name: "Customer", Email: "customer@example.com"}
|
|
require.NoError(t, db.Create(contact).Error)
|
|
conversation := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusPending), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType}
|
|
require.NoError(t, db.Create(conversation).Error)
|
|
assistant := &model.CaptainAssistant{AccountID: account.ID, Name: "Fin", Config: []byte(`{"handoff_message":"Let me connect you."}`), Status: model.AssistantStatusActive}
|
|
require.NoError(t, db.Create(assistant).Error)
|
|
require.NoError(t, db.Create(&model.CaptainInbox{AccountID: account.ID, AssistantID: assistant.ID, InboxID: inbox.ID}).Error)
|
|
bot := &model.AgentBot{AccountID: &account.ID, Name: "Captain", BotType: "captain", Config: []byte(fmt.Sprintf(`{"assistant_id":%d}`, assistant.ID))}
|
|
require.NoError(t, db.Create(bot).Error)
|
|
require.NoError(t, db.Create(&model.AgentBotInbox{AgentBotID: bot.ID, InboxID: inbox.ID, Status: model.AgentBotInboxActive}).Error)
|
|
require.NoError(t, db.Model(conversation).Updates(map[string]any{"assignee_agent_bot_id": bot.ID, "ai_takeover_version": 1}).Error)
|
|
conversation.AssigneeAgentBotID = &bot.ID
|
|
conversation.AITakeoverVersion = 1
|
|
conversationSvc := NewCaptainConversationService(db, nil)
|
|
messageSvc := NewMessageService(repository.NewMessageRepo(db), channel.NewDispatcher(), nil)
|
|
return db, conversationSvc, messageSvc, account, inbox, conversation, assistant
|
|
}
|
|
|
|
func TestCaptainConversationResponseJobQueuesFromIncomingMessage(t *testing.T) {
|
|
db, conversationSvc, messageSvc, account, _, conversation, assistant := setupCaptainConversationWorkerTest(t)
|
|
conversationSvc.SetResponseBackend(&fakeCaptainConversationBackend{response: &CaptainConversationResponse{Content: "Welcome to Captain", AgentName: "Fin"}})
|
|
wp := worker.NewWorkerPool(db)
|
|
conversationSvc.SetWorkerPool(wp)
|
|
messageSvc.SetWorkerPool(wp)
|
|
|
|
incoming, err := messageSvc.Create(context.Background(), account.ID, 99, CreateMessageRequest{ConversationID: conversation.ID, Content: "Hello", MessageType: string(model.MessageTypeIncoming), ContentType: string(model.MessageContentTypeText)})
|
|
require.NoError(t, err)
|
|
require.NotZero(t, incoming.ID)
|
|
|
|
var count int64
|
|
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ? AND status = ?", TaskTypeCaptainConversationResponseBuilder, model.BackgroundJobStatusQueued).Count(&count).Error)
|
|
assert.Equal(t, int64(1), count)
|
|
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
require.NoError(t, err)
|
|
assert.True(t, processed)
|
|
|
|
var outgoing model.Message
|
|
require.NoError(t, db.Where("conversation_id = ? AND message_type = ?", conversation.ID, model.MessageTypeOutgoing).First(&outgoing).Error)
|
|
assert.Equal(t, assistant.ID, *outgoing.SenderID)
|
|
assert.Equal(t, "Captain::Assistant", outgoing.SenderType)
|
|
assert.Equal(t, "Welcome to Captain", outgoing.Content)
|
|
|
|
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeMessageSendReply).Count(&count).Error)
|
|
assert.Equal(t, int64(1), count)
|
|
}
|
|
|
|
func TestCaptainConversationResponseUsesShangwutongDelivery(t *testing.T) {
|
|
db, conversationSvc, messageSvc, account, inbox, conversation, _ := setupCaptainConversationWorkerTest(t)
|
|
require.NoError(t, db.Model(&model.Inbox{}).Where("id = ?", inbox.ID).Update("channel_type", "shangwutong").Error)
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).Updates(map[string]any{"channel_type": "shangwutong", "channel": "shangwutong"}).Error)
|
|
conversationSvc.SetResponseBackend(&fakeCaptainConversationBackend{response: &CaptainConversationResponse{Content: "AI reply"}})
|
|
wp := worker.NewWorkerPool(db)
|
|
conversationSvc.SetWorkerPool(wp)
|
|
messageSvc.SetWorkerPool(wp)
|
|
conversationSvc.SetMessageService(messageSvc)
|
|
|
|
var bot model.AgentBot
|
|
require.NoError(t, db.Where("account_id = ?", account.ID).FirstOrCreate(&bot, model.AgentBot{AccountID: &account.ID, Name: "Captain", BotType: "captain", Config: []byte(`{"assistant_id":1}`)}).Error)
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).Update("assignee_agent_bot_id", bot.ID).Error)
|
|
message, err := conversationSvc.BuildConversationResponseByAccount(context.Background(), account.ID, conversation.ID, 1)
|
|
require.NoError(t, err)
|
|
require.NotNil(t, message)
|
|
assert.Equal(t, string(model.MessageStatusProgress), message.Status)
|
|
|
|
var count int64
|
|
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeShangwutongWebhookDelivery).Count(&count).Error)
|
|
assert.Equal(t, int64(1), count)
|
|
}
|
|
|
|
func TestCaptainConversationResponseJobHandoffOpensConversation(t *testing.T) {
|
|
db, conversationSvc, _, account, _, conversation, _ := setupCaptainConversationWorkerTest(t)
|
|
conversationSvc.SetResponseBackend(&fakeCaptainConversationBackend{response: &CaptainConversationResponse{Action: "handoff"}})
|
|
wp := worker.NewWorkerPool(db)
|
|
conversationSvc.SetWorkerPool(wp)
|
|
|
|
_, err := wp.Enqueue(context.Background(), TaskTypeCaptainConversationResponseBuilder, captainConversationResponseBuilderJob{AccountID: account.ID, ConversationID: conversation.ID, AssistantID: 1, MessageID: 1}, worker.WithMaxAttempts(3))
|
|
require.NoError(t, err)
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
require.NoError(t, err)
|
|
assert.True(t, processed)
|
|
|
|
var updated model.Conversation
|
|
require.NoError(t, db.First(&updated, conversation.ID).Error)
|
|
assert.Equal(t, string(model.ConversationStatusOpen), updated.Status)
|
|
var outgoing model.Message
|
|
require.NoError(t, db.Where("conversation_id = ? AND message_type = ?", conversation.ID, model.MessageTypeOutgoing).First(&outgoing).Error)
|
|
assert.Equal(t, "Let me connect you.", outgoing.Content)
|
|
}
|
|
|
|
func TestCaptainConversationResponseJobRetriesWhenProviderDisabled(t *testing.T) {
|
|
db, conversationSvc, _, account, _, conversation, _ := setupCaptainConversationWorkerTest(t)
|
|
now := time.Date(2026, 6, 6, 4, 30, 0, 0, time.UTC)
|
|
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }), worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }))
|
|
conversationSvc.SetWorkerPool(wp)
|
|
|
|
_, err := wp.Enqueue(context.Background(), TaskTypeCaptainConversationResponseBuilder, captainConversationResponseBuilderJob{AccountID: account.ID, ConversationID: conversation.ID, AssistantID: 1, MessageID: 1}, worker.WithMaxAttempts(3))
|
|
require.NoError(t, err)
|
|
processed, err := wp.ProcessOne(context.Background())
|
|
require.Error(t, err)
|
|
assert.True(t, processed)
|
|
|
|
var job model.BackgroundJob
|
|
require.NoError(t, db.Where("job_type = ?", TaskTypeCaptainConversationResponseBuilder).First(&job).Error)
|
|
assert.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
|
|
assert.Contains(t, job.LastError, "captain conversation response generation disabled")
|
|
}
|
|
|
|
func TestCaptainConversationResponseSkipsNonPendingConversation(t *testing.T) {
|
|
db, conversationSvc, messageSvc, account, _, conversation, _ := setupCaptainConversationWorkerTest(t)
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).Update("status", string(model.ConversationStatusOpen)).Error)
|
|
wp := worker.NewWorkerPool(db)
|
|
conversationSvc.SetWorkerPool(wp)
|
|
messageSvc.SetWorkerPool(wp)
|
|
|
|
_, err := messageSvc.Create(context.Background(), account.ID, 99, CreateMessageRequest{ConversationID: conversation.ID, Content: "Hello", MessageType: string(model.MessageTypeIncoming), ContentType: string(model.MessageContentTypeText)})
|
|
require.NoError(t, err)
|
|
var count int64
|
|
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeCaptainConversationResponseBuilder).Count(&count).Error)
|
|
assert.Equal(t, int64(0), count)
|
|
}
|
|
|
|
func TestCaptainInboxAutoResponseIsOffByDefault(t *testing.T) {
|
|
db, _, messageSvc, account, _, conversation, _ := setupCaptainConversationWorkerTest(t)
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).Update("assignee_agent_bot_id", nil).Error)
|
|
wp := worker.NewWorkerPool(db)
|
|
messageSvc.SetWorkerPool(wp)
|
|
|
|
_, err := messageSvc.Create(context.Background(), account.ID, 99, CreateMessageRequest{ConversationID: conversation.ID, Content: "Hello", MessageType: string(model.MessageTypeIncoming), ContentType: string(model.MessageContentTypeText)})
|
|
require.NoError(t, err)
|
|
|
|
var count int64
|
|
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeCaptainConversationResponseBuilder).Count(&count).Error)
|
|
assert.Zero(t, count)
|
|
}
|
|
|
|
type fakeCaptainConversationBackend struct {
|
|
response *CaptainConversationResponse
|
|
err error
|
|
beforeReturn func()
|
|
}
|
|
|
|
func (b *fakeCaptainConversationBackend) GenerateCaptainConversationResponse(ctx context.Context, req CaptainConversationResponseRequest) (*CaptainConversationResponse, error) {
|
|
if b.err != nil {
|
|
return nil, b.err
|
|
}
|
|
if b.response == nil {
|
|
return nil, errors.New("missing response")
|
|
}
|
|
if b.beforeReturn != nil {
|
|
b.beforeReturn()
|
|
}
|
|
return b.response, nil
|
|
}
|
|
|
|
func TestCaptainConversationResponseRejectsStaleJobAfterReentry(t *testing.T) {
|
|
db, conversationSvc, messageSvc, account, _, conversation, assistant := setupCaptainConversationWorkerTest(t)
|
|
conversationSvc.SetMessageService(messageSvc)
|
|
conversationSvc.SetResponseBackend(&fakeCaptainConversationBackend{response: &CaptainConversationResponse{Content: "stale"}})
|
|
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).
|
|
Updates(map[string]any{"ai_takeover_version": 3}).Error)
|
|
message, err := conversationSvc.buildConversationResponseByAccount(context.Background(), account.ID, conversation.ID, assistant.ID, 1)
|
|
require.NoError(t, err)
|
|
assert.Nil(t, message)
|
|
|
|
var count int64
|
|
require.NoError(t, db.Model(&model.Message{}).Where("conversation_id = ? AND message_type = ?", conversation.ID, model.MessageTypeOutgoing).Count(&count).Error)
|
|
assert.Zero(t, count)
|
|
}
|
|
|
|
func TestCaptainConversationResponseRechecksTakeoverBeforeCreate(t *testing.T) {
|
|
db, conversationSvc, messageSvc, account, _, conversation, assistant := setupCaptainConversationWorkerTest(t)
|
|
conversationSvc.SetMessageService(messageSvc)
|
|
conversationSvc.SetResponseBackend(&fakeCaptainConversationBackend{
|
|
response: &CaptainConversationResponse{Content: "late"},
|
|
beforeReturn: func() {
|
|
require.NoError(t, db.Model(&model.Conversation{}).Where("id = ?", conversation.ID).
|
|
Updates(map[string]any{"assignee_agent_bot_id": nil, "status": model.ConversationStatusOpen, "ai_takeover_version": 2}).Error)
|
|
},
|
|
})
|
|
message, err := conversationSvc.buildConversationResponseByAccount(context.Background(), account.ID, conversation.ID, assistant.ID, 1)
|
|
require.NoError(t, err)
|
|
assert.Nil(t, message)
|
|
}
|