From 0540ae3791687d50089010254b301f9578ad0d3a Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 13 Aug 2026 16:24:39 +0800 Subject: [PATCH] H-55: make Captain bindings atomic (#9) Co-authored-by: Rogee --- backend/internal/model/agent_bot.go | 27 +++-- backend/internal/model/agent_bot_inbox.go | 16 +-- backend/internal/model/captain_models.go | 2 +- .../internal/repository/conversation_repo.go | 9 +- backend/internal/service/ai_takeover_test.go | 7 +- .../service/captain_assistant_service.go | 91 ++++++++------ .../service/captain_binding_atomicity_test.go | 112 ++++++++++++++++++ backend/internal/service/widget_service.go | 83 ++++++------- ...0079_make_captain_bindings_unique.down.sql | 4 + ...000079_make_captain_bindings_unique.up.sql | 99 ++++++++++++++++ 10 files changed, 346 insertions(+), 104 deletions(-) create mode 100644 backend/internal/service/captain_binding_atomicity_test.go create mode 100644 backend/migrations/000079_make_captain_bindings_unique.down.sql create mode 100644 backend/migrations/000079_make_captain_bindings_unique.up.sql diff --git a/backend/internal/model/agent_bot.go b/backend/internal/model/agent_bot.go index 0b2cc333..4fde2c57 100644 --- a/backend/internal/model/agent_bot.go +++ b/backend/internal/model/agent_bot.go @@ -11,18 +11,19 @@ import ( // Bot can be bound to Inbox via AgentBotInbox, or assigned as conversation assignee. // Account-level bot (account_id != nil) vs global bot (account_id == nil). type AgentBot struct { - ID uint `gorm:"primaryKey" json:"id"` - AccountID *uint `gorm:"index" json:"account_id,omitempty"` - Name string `gorm:"size:255;not null" json:"name"` - Description string `gorm:"size:512" json:"description"` - AvatarURL string `gorm:"size:512" json:"avatar_url"` - OutgoingURL string `gorm:"size:1024" json:"outgoing_url"` // Webhook push URL - BotType string `gorm:"size:50;default:'webhook'" json:"bot_type"` // webhook/default/custom - Secret string `gorm:"size:128;uniqueIndex" json:"secret,omitempty"` // Webhook signing secret (generated on create/reset) - AccessToken string `gorm:"size:128;uniqueIndex" json:"access_token,omitempty"` // Bot API access token (generated on create/reset) - Config json.RawMessage `gorm:"type:jsonb" json:"config"` - CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` - UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` + ID uint `gorm:"primaryKey" json:"id"` + AccountID *uint `gorm:"index;uniqueIndex:idx_agent_bots_captain_assistant" json:"account_id,omitempty"` + CaptainAssistantID *uint `gorm:"uniqueIndex:idx_agent_bots_captain_assistant" json:"-"` + Name string `gorm:"size:255;not null" json:"name"` + Description string `gorm:"size:512" json:"description"` + AvatarURL string `gorm:"size:512" json:"avatar_url"` + OutgoingURL string `gorm:"size:1024" json:"outgoing_url"` // Webhook push URL + BotType string `gorm:"size:50;default:'webhook'" json:"bot_type"` // webhook/default/custom + Secret string `gorm:"size:128;uniqueIndex" json:"secret,omitempty"` // Webhook signing secret (generated on create/reset) + AccessToken string `gorm:"size:128;uniqueIndex" json:"access_token,omitempty"` // Bot API access token (generated on create/reset) + Config json.RawMessage `gorm:"type:jsonb" json:"config"` + CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` + UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` Inboxes []AgentBotInbox `gorm:"foreignKey:AgentBotID" json:"inboxes,omitempty"` } @@ -38,4 +39,4 @@ func (b *AgentBot) PushEventData() map[string]interface{} { } } -func (AgentBot) TableName() string { return "agent_bots" } \ No newline at end of file +func (AgentBot) TableName() string { return "agent_bots" } diff --git a/backend/internal/model/agent_bot_inbox.go b/backend/internal/model/agent_bot_inbox.go index eff982e8..cd3ecc7d 100644 --- a/backend/internal/model/agent_bot_inbox.go +++ b/backend/internal/model/agent_bot_inbox.go @@ -17,13 +17,13 @@ const ( // Supports active/inactive status control — allows pausing bot monitoring without deleting binding. // account_id is auto-populated from inbox.account_id (ensure_account_id callback). type AgentBotInbox struct { - ID uint `gorm:"primaryKey" json:"id"` - AgentBotID uint `gorm:"not null;index" json:"agent_bot_id"` - InboxID uint `gorm:"not null;index" json:"inbox_id"` - AccountID *uint `gorm:"index" json:"account_id,omitempty"` // inherited from Inbox - Status AgentBotInboxStatus `gorm:"default:0" json:"status"` // 0=active, 1=inactive - CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` - UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` + ID uint `gorm:"primaryKey" json:"id"` + AgentBotID uint `gorm:"not null;index;uniqueIndex:idx_agent_bot_inboxes_bot_inbox" json:"agent_bot_id"` + InboxID uint `gorm:"not null;index;uniqueIndex:idx_agent_bot_inboxes_bot_inbox" json:"inbox_id"` + AccountID *uint `gorm:"index" json:"account_id,omitempty"` // inherited from Inbox + Status AgentBotInboxStatus `gorm:"default:0" json:"status"` // 0=active, 1=inactive + CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` + UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` AgentBot AgentBot `gorm:"foreignKey:AgentBotID" json:"agent_bot,omitempty"` Inbox Inbox `gorm:"foreignKey:InboxID" json:"inbox,omitempty"` @@ -34,4 +34,4 @@ func (abi *AgentBotInbox) IsActive() bool { return abi.Status == AgentBotInboxActive } -func (AgentBotInbox) TableName() string { return "agent_bot_inboxes" } \ No newline at end of file +func (AgentBotInbox) TableName() string { return "agent_bot_inboxes" } diff --git a/backend/internal/model/captain_models.go b/backend/internal/model/captain_models.go index f852b782..77fab511 100644 --- a/backend/internal/model/captain_models.go +++ b/backend/internal/model/captain_models.go @@ -291,7 +291,7 @@ func (CaptainCustomTool) TableName() string { return "captain_custom_tools" } type CaptainInbox struct { Base AssistantID uint `gorm:"column:captain_assistant_id;index;not null" json:"captain_assistant_id"` - InboxID uint `gorm:"index;not null" json:"inbox_id"` + InboxID uint `gorm:"index;not null;uniqueIndex:idx_captain_inboxes_active_inbox,where:deleted_at IS NULL" json:"inbox_id"` AccountID uint `gorm:"index;not null" json:"account_id"` } diff --git a/backend/internal/repository/conversation_repo.go b/backend/internal/repository/conversation_repo.go index 8147840c..3bc316cf 100644 --- a/backend/internal/repository/conversation_repo.go +++ b/backend/internal/repository/conversation_repo.go @@ -272,11 +272,16 @@ func (r *ConversationRepo) Search(ctx context.Context, accountID uint, query str // Create inserts a new conversation. func (r *ConversationRepo) Create(ctx context.Context, conversation *model.Conversation) error { + return r.CreateWithDB(ctx, r.db, conversation) +} + +// CreateWithDB inserts a conversation using the caller's transaction. +func (r *ConversationRepo) CreateWithDB(ctx context.Context, db *gorm.DB, conversation *model.Conversation) error { if conversation.DisplayID != nil && *conversation.DisplayID != 0 { - return r.db.WithContext(ctx).Create(conversation).Error + return db.WithContext(ctx).Create(conversation).Error } - return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { // PostgreSQL needs serialization because MAX(display_id)+1 is otherwise // racy when a channel imports several conversations concurrently. if tx.Dialector != nil && tx.Dialector.Name() == "postgres" { diff --git a/backend/internal/service/ai_takeover_test.go b/backend/internal/service/ai_takeover_test.go index 8e7403ce..fdc9af84 100644 --- a/backend/internal/service/ai_takeover_test.go +++ b/backend/internal/service/ai_takeover_test.go @@ -79,8 +79,11 @@ func TestConversationServiceAITakeoverRejectsMissingOrAmbiguousAI(t *testing.T) _, err := svc.StartAITakeover(context.Background(), account.ID, conversation.ID) assert.ErrorContains(t, err, "exactly one active AI") - configureInboxAI(t, svc, account.ID, inbox.ID) - configureInboxAI(t, svc, account.ID, inbox.ID) + bot, _ := configureInboxAI(t, svc, account.ID, inbox.ID) + otherBot := &model.AgentBot{AccountID: &account.ID, Name: "Other Captain", BotType: "captain", Secret: "other-secret", AccessToken: "other-token", Config: []byte(`{"assistant_id":999}`)} + require.NoError(t, db.Create(otherBot).Error) + require.NoError(t, db.Create(&model.AgentBotInbox{AgentBotID: otherBot.ID, InboxID: inbox.ID, Status: model.AgentBotInboxActive}).Error) + require.NotEqual(t, bot.ID, otherBot.ID) _, err = svc.StartAITakeover(context.Background(), account.ID, conversation.ID) assert.ErrorContains(t, err, "exactly one active AI") } diff --git a/backend/internal/service/captain_assistant_service.go b/backend/internal/service/captain_assistant_service.go index ea169082..f7414f57 100644 --- a/backend/internal/service/captain_assistant_service.go +++ b/backend/internal/service/captain_assistant_service.go @@ -17,6 +17,7 @@ import ( "github.com/pgvector/pgvector-go" "github.com/redis/go-redis/v9" "gorm.io/gorm" + "gorm.io/gorm/clause" ) // CaptainAssistantService implements business logic for CaptainAssistant operations. @@ -571,9 +572,20 @@ func (s *CaptainAssistantService) AssociateInbox(ctx context.Context, assistantI db := s.assistantRepo.DB().WithContext(ctx) err = db.Transaction(func(tx *gorm.DB) error { ci := &model.CaptainInbox{AssistantID: assistantID, InboxID: inboxID, AccountID: accountID} - if err := tx.Where("captain_assistant_id = ? AND inbox_id = ?", assistantID, inboxID).FirstOrCreate(ci).Error; err != nil { + if err := tx.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "inbox_id"}}, + TargetWhere: clause.Where{Exprs: []clause.Expression{clause.Expr{SQL: "deleted_at IS NULL"}}}, + DoNothing: true, + }).Create(ci).Error; err != nil { return err } + var persisted model.CaptainInbox + if err := tx.Where("inbox_id = ?", inboxID).First(&persisted).Error; err != nil { + return err + } + if persisted.AssistantID != assistantID { + return errors.New("inbox is already associated with another assistant") + } _, err := ensureCaptainAgentBotBinding(ctx, tx, assistant, inboxID) return err @@ -595,52 +607,65 @@ func (s *CaptainAssistantService) DissociateInbox(ctx context.Context, accountID } else if ci.AccountID != accountID { return fmt.Errorf("captain inbox not found") } - if err := s.inboxRepo.DeleteByAccount(ctx, accountID, assistantID, inboxID); err != nil { + db := s.assistantRepo.DB().WithContext(ctx) + if err := db.Transaction(func(tx *gorm.DB) error { + if err := tx.Where("account_id = ? AND captain_assistant_id = ? AND inbox_id = ?", accountID, assistantID, inboxID).Delete(&model.CaptainInbox{}).Error; err != nil { + return err + } + var bot model.AgentBot + if err := tx.Where("account_id = ? AND captain_assistant_id = ?", accountID, assistantID).First(&bot).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil + } + return err + } + return tx.Where("inbox_id = ? AND agent_bot_id = ?", inboxID, bot.ID).Delete(&model.AgentBotInbox{}).Error + }); err != nil { applogger.L().Errorf("DissociateInbox: %v", err) return fmt.Errorf("dissociate inbox: %w", err) } - var bots []model.AgentBot - if err := s.assistantRepo.DB().WithContext(ctx).Where("account_id = ? AND bot_type = ?", accountID, "captain").Find(&bots).Error; err != nil { - return err - } - for _, bot := range bots { - if extractAssistantIDFromBotConfig(bot.Config) == assistantID { - if err := s.assistantRepo.DB().WithContext(ctx).Where("inbox_id = ? AND agent_bot_id = ?", inboxID, bot.ID).Delete(&model.AgentBotInbox{}).Error; err != nil { - return fmt.Errorf("dissociate inbox bot: %w", err) - } - } - } return nil } func ensureCaptainAgentBotBinding(ctx context.Context, db *gorm.DB, assistant *model.CaptainAssistant, inboxID uint) (*model.AgentBot, error) { - var bots []model.AgentBot - if err := db.WithContext(ctx).Where("account_id = ? AND bot_type = ?", assistant.AccountID, "captain").Find(&bots).Error; err != nil { + var legacyBots []model.AgentBot + if err := db.WithContext(ctx).Where("account_id = ? AND bot_type = ? AND captain_assistant_id IS NULL", assistant.AccountID, "captain").Find(&legacyBots).Error; err != nil { return nil, err } - var bot *model.AgentBot - for i := range bots { - if extractAssistantIDFromBotConfig(bots[i].Config) == assistant.ID { - bot = &bots[i] + for i := range legacyBots { + if extractAssistantIDFromBotConfig(legacyBots[i].Config) == assistant.ID { + if err := db.WithContext(ctx).Model(&legacyBots[i]).Where("captain_assistant_id IS NULL").Update("captain_assistant_id", assistant.ID).Error; err != nil { + return nil, err + } break } } - if bot == nil { - token, err := generateBotAccessToken() - if err != nil { - return nil, err - } - secret, err := generateBotSecret() - if err != nil { - return nil, err - } - bot = &model.AgentBot{AccountID: &assistant.AccountID, Name: fmt.Sprintf("Captain Assistant #%d", assistant.ID), Description: assistant.Name, BotType: "captain", AccessToken: token, Secret: secret, Config: json.RawMessage(fmt.Sprintf(`{"assistant_id":%d}`, assistant.ID))} - if err := db.WithContext(ctx).Create(bot).Error; err != nil { - return nil, err - } + + token, err := generateBotAccessToken() + if err != nil { + return nil, err } + secret, err := generateBotSecret() + if err != nil { + return nil, err + } + bot := &model.AgentBot{AccountID: &assistant.AccountID, CaptainAssistantID: &assistant.ID, Name: fmt.Sprintf("Captain Assistant #%d", assistant.ID), Description: assistant.Name, BotType: "captain", AccessToken: token, Secret: secret, Config: json.RawMessage(fmt.Sprintf(`{"assistant_id":%d}`, assistant.ID))} + if err := db.WithContext(ctx).Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "account_id"}, {Name: "captain_assistant_id"}}, + DoNothing: true, + }).Create(bot).Error; err != nil { + return nil, err + } + var persisted model.AgentBot + if err := db.WithContext(ctx).Where("account_id = ? AND captain_assistant_id = ?", assistant.AccountID, assistant.ID).First(&persisted).Error; err != nil { + return nil, fmt.Errorf("load captain bot: %w", err) + } + bot = &persisted binding := &model.AgentBotInbox{AgentBotID: bot.ID, InboxID: inboxID, AccountID: &assistant.AccountID, Status: model.AgentBotInboxActive} - if err := db.WithContext(ctx).Where("agent_bot_id = ? AND inbox_id = ?", bot.ID, inboxID).Assign("status", model.AgentBotInboxActive).FirstOrCreate(binding).Error; err != nil { + if err := db.WithContext(ctx).Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "agent_bot_id"}, {Name: "inbox_id"}}, + DoUpdates: clause.Assignments(map[string]any{"status": model.AgentBotInboxActive, "account_id": assistant.AccountID}), + }).Create(binding).Error; err != nil { return nil, err } return bot, nil diff --git a/backend/internal/service/captain_binding_atomicity_test.go b/backend/internal/service/captain_binding_atomicity_test.go new file mode 100644 index 00000000..4fc69ddf --- /dev/null +++ b/backend/internal/service/captain_binding_atomicity_test.go @@ -0,0 +1,112 @@ +package service + +import ( + "context" + "errors" + "sync" + "testing" + + "github.com/gochat/gochat/internal/model" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "gorm.io/gorm/logger" +) + +func TestWidgetConversationRollsBackCaptainBindingOnCreateFailure(t *testing.T) { + db, svc := setupWidgetServiceTest(t) + account, inbox := seedWidgetInbox(t, db) + contact := &model.Contact{AccountID: account.ID, Name: "Atomic visitor"} + require.NoError(t, db.Create(contact).Error) + contactInbox := &model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: "atomic"} + require.NoError(t, db.Create(contactInbox).Error) + assistant := &model.CaptainAssistant{AccountID: account.ID, Name: "Atomic", Status: model.AssistantStatusActive} + require.NoError(t, db.Create(assistant).Error) + require.NoError(t, db.Create(&model.CaptainInbox{AccountID: account.ID, InboxID: inbox.ID, AssistantID: assistant.ID}).Error) + require.NoError(t, db.Create(&model.CaptainPreference{AccountID: account.ID, AutoReplyEnabled: true}).Error) + require.NoError(t, db.Callback().Create().Before("gorm:create").Register("test:fail_conversation_create", func(tx *gorm.DB) { + if tx.Statement.Table == "conversations" { + tx.AddError(errors.New("conversation create failed")) + } + })) + t.Cleanup(func() { _ = db.Callback().Create().Remove("test:fail_conversation_create") }) + + _, err := svc.createWidgetConversation(context.Background(), contactInbox, nil, nil) + require.ErrorContains(t, err, "conversation create failed") + for _, value := range []any{&model.Conversation{}, &model.AgentBot{}, &model.AgentBotInbox{}} { + var count int64 + require.NoError(t, db.Model(value).Count(&count).Error) + assert.Zero(t, count) + } +} + +func TestEnsureCaptainAgentBotBindingConcurrentCallsStayUnique(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) + require.NoError(t, err) + require.NoError(t, db.AutoMigrate(&model.CaptainAssistant{}, &model.AgentBot{}, &model.AgentBotInbox{})) + sqlDB, err := db.DB() + require.NoError(t, err) + sqlDB.SetMaxOpenConns(1) + t.Cleanup(func() { _ = sqlDB.Close() }) + assistant := &model.CaptainAssistant{AccountID: 1, Name: "Concurrent", Status: model.AssistantStatusActive} + require.NoError(t, db.Create(assistant).Error) + + const calls = 8 + errs := make(chan error, calls) + ids := make(chan uint, calls) + var wg sync.WaitGroup + for range calls { + wg.Add(1) + go func() { + defer wg.Done() + bot, err := ensureCaptainAgentBotBinding(context.Background(), db, assistant, 42) + errs <- err + if bot != nil { + ids <- bot.ID + } + }() + } + wg.Wait() + close(errs) + close(ids) + for err := range errs { + require.NoError(t, err) + } + var first uint + for id := range ids { + if first == 0 { + first = id + } + assert.Equal(t, first, id) + } + for _, value := range []any{&model.AgentBot{}, &model.AgentBotInbox{}} { + var count int64 + require.NoError(t, db.Model(value).Count(&count).Error) + assert.Equal(t, int64(1), count) + } +} + +func TestDissociateInboxRollsBackCaptainDeleteWhenBindingDeleteFails(t *testing.T) { + db, _, svc := setupCov7CaptainAssistant(t) + account := seedCov7Account(t, db) + inbox := seedCov7Inbox(t, db, account.ID) + assistant, err := svc.Create(context.Background(), account.ID, &CreateAssistantRequest{Name: "Rollback", Description: "test"}) + require.NoError(t, err) + _, err = svc.AssociateInbox(context.Background(), assistant.ID, inbox.ID, account.ID) + require.NoError(t, err) + require.NoError(t, db.Callback().Delete().Before("gorm:delete").Register("test:fail_binding_delete", func(tx *gorm.DB) { + if tx.Statement.Table == "agent_bot_inboxes" { + tx.AddError(errors.New("binding delete failed")) + } + })) + t.Cleanup(func() { _ = db.Callback().Delete().Remove("test:fail_binding_delete") }) + + err = svc.DissociateInbox(context.Background(), account.ID, assistant.ID, inbox.ID) + require.ErrorContains(t, err, "binding delete failed") + var captainInboxCount, bindingCount int64 + require.NoError(t, db.Model(&model.CaptainInbox{}).Where("inbox_id = ?", inbox.ID).Count(&captainInboxCount).Error) + require.NoError(t, db.Model(&model.AgentBotInbox{}).Where("inbox_id = ?", inbox.ID).Count(&bindingCount).Error) + assert.Equal(t, int64(1), captainInboxCount) + assert.Equal(t, int64(1), bindingCount) +} diff --git a/backend/internal/service/widget_service.go b/backend/internal/service/widget_service.go index 879ec581..fbc1715a 100644 --- a/backend/internal/service/widget_service.go +++ b/backend/internal/service/widget_service.go @@ -1720,57 +1720,50 @@ func (s *WidgetService) createWidgetConversation(ctx context.Context, contactInb return nil, err } - status := string(model.ConversationStatusOpen) - var agentBotID *uint - var takeoverVersion uint db := s.conversationRepo.DB().WithContext(ctx) - if db.Migrator().HasTable(&model.CaptainPreference{}) && db.Migrator().HasTable(&model.CaptainInbox{}) { - var autoReplyCount int64 - if err := db.Model(&model.CaptainPreference{}). - Where("account_id = ? AND auto_reply_enabled = ?", inbox.AccountID, true).Count(&autoReplyCount).Error; err != nil { - return nil, err - } - if autoReplyCount > 0 { - var captainInbox model.CaptainInbox - if err := db.Where("account_id = ? AND inbox_id = ?", inbox.AccountID, inbox.ID).First(&captainInbox).Error; err != nil { - if errors.Is(err, gorm.ErrRecordNotFound) { - goto create + conversation := &model.Conversation{ + AccountID: inbox.AccountID, + InboxID: inbox.ID, + ContactID: contactInbox.ContactID, + ContactInboxID: &contactInbox.ID, + Status: string(model.ConversationStatusOpen), + ChannelType: inbox.ChannelType, + Channel: inbox.ChannelType, + CustomAttributes: mustJSON(customAttributes), + Labels: strings.Join(s.validWidgetLabels(ctx, inbox.AccountID, labels), ","), + } + if err := db.Transaction(func(tx *gorm.DB) error { + if tx.Migrator().HasTable(&model.CaptainPreference{}) && tx.Migrator().HasTable(&model.CaptainInbox{}) { + var autoReplyCount int64 + if err := tx.Model(&model.CaptainPreference{}).Where("account_id = ? AND auto_reply_enabled = ?", inbox.AccountID, true).Count(&autoReplyCount).Error; err != nil { + return err + } + if autoReplyCount > 0 { + var captainInbox model.CaptainInbox + if err := tx.Where("account_id = ? AND inbox_id = ?", inbox.AccountID, inbox.ID).First(&captainInbox).Error; err != nil { + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + } else { + var assistant model.CaptainAssistant + if err := tx.Where("account_id = ? AND id = ? AND status = ?", inbox.AccountID, captainInbox.AssistantID, model.AssistantStatusActive).First(&assistant).Error; err != nil { + return err + } + bot, err := ensureCaptainAgentBotBinding(ctx, tx, &assistant, inbox.ID) + if err != nil { + return err + } + conversation.Status = string(model.ConversationStatusPending) + conversation.AssigneeAgentBotID = &bot.ID + conversation.AITakeoverVersion = 1 } - return nil, err } - var assistant model.CaptainAssistant - if err := db.Where("account_id = ? AND id = ? AND status = ?", inbox.AccountID, captainInbox.AssistantID, model.AssistantStatusActive).First(&assistant).Error; err != nil { - return nil, err - } - bot, err := ensureCaptainAgentBotBinding(ctx, db, &assistant, inbox.ID) - if err != nil { - return nil, err - } - status = string(model.ConversationStatusPending) - agentBotID = &bot.ID - takeoverVersion = 1 } - } - -create: - conversation := model.Conversation{ - AccountID: inbox.AccountID, - InboxID: inbox.ID, - ContactID: contactInbox.ContactID, - ContactInboxID: &contactInbox.ID, - Status: status, - AssigneeAgentBotID: agentBotID, - AITakeoverVersion: takeoverVersion, - ChannelType: inbox.ChannelType, - Channel: inbox.ChannelType, - CustomAttributes: mustJSON(customAttributes), - Labels: strings.Join(s.validWidgetLabels(ctx, inbox.AccountID, labels), ","), - } - - if err := s.conversationRepo.Create(ctx, &conversation); err != nil { + return s.conversationRepo.CreateWithDB(ctx, tx, conversation) + }); err != nil { return nil, err } - return &conversation, nil + return conversation, nil } func (s *WidgetService) validWidgetLabels(ctx context.Context, accountID uint, labels []string) []string { diff --git a/backend/migrations/000079_make_captain_bindings_unique.down.sql b/backend/migrations/000079_make_captain_bindings_unique.down.sql new file mode 100644 index 00000000..f74b2e12 --- /dev/null +++ b/backend/migrations/000079_make_captain_bindings_unique.down.sql @@ -0,0 +1,4 @@ +DROP INDEX IF EXISTS idx_captain_inboxes_active_inbox; +DROP INDEX IF EXISTS idx_agent_bot_inboxes_bot_inbox; +DROP INDEX IF EXISTS idx_agent_bots_captain_assistant; +ALTER TABLE agent_bots DROP COLUMN IF EXISTS captain_assistant_id; diff --git a/backend/migrations/000079_make_captain_bindings_unique.up.sql b/backend/migrations/000079_make_captain_bindings_unique.up.sql new file mode 100644 index 00000000..ed0e12dc --- /dev/null +++ b/backend/migrations/000079_make_captain_bindings_unique.up.sql @@ -0,0 +1,99 @@ +ALTER TABLE agent_bots ADD COLUMN IF NOT EXISTS captain_assistant_id BIGINT; + +UPDATE agent_bots +SET captain_assistant_id = (config ->> 'assistant_id')::BIGINT +WHERE bot_type = 'captain' + AND captain_assistant_id IS NULL + AND config ->> 'assistant_id' ~ '^[0-9]+$'; + +UPDATE conversations conversation +SET assignee_agent_bot_id = canonical.keep_id +FROM ( + SELECT account_id, captain_assistant_id, MIN(id) AS keep_id, + ARRAY_AGG(id) AS duplicate_ids + FROM agent_bots + WHERE captain_assistant_id IS NOT NULL + GROUP BY account_id, captain_assistant_id + HAVING COUNT(*) > 1 +) canonical +WHERE conversation.assignee_agent_bot_id = ANY(canonical.duplicate_ids) + AND conversation.assignee_agent_bot_id <> canonical.keep_id; + +UPDATE agent_bot_inboxes binding +SET agent_bot_id = canonical.keep_id +FROM ( + SELECT account_id, captain_assistant_id, MIN(id) AS keep_id, + ARRAY_AGG(id) AS duplicate_ids + FROM agent_bots + WHERE captain_assistant_id IS NOT NULL + GROUP BY account_id, captain_assistant_id + HAVING COUNT(*) > 1 +) canonical +WHERE binding.agent_bot_id = ANY(canonical.duplicate_ids) + AND binding.agent_bot_id <> canonical.keep_id; + +UPDATE agent_bot_presence_events event +SET agent_bot_id = canonical.keep_id +FROM ( + SELECT account_id, captain_assistant_id, MIN(id) AS keep_id, + ARRAY_AGG(id) AS duplicate_ids + FROM agent_bots + WHERE captain_assistant_id IS NOT NULL + GROUP BY account_id, captain_assistant_id + HAVING COUNT(*) > 1 +) canonical +WHERE event.agent_bot_id = ANY(canonical.duplicate_ids) + AND event.agent_bot_id <> canonical.keep_id; + +UPDATE bot_rules rule +SET agent_bot_id = canonical.keep_id +FROM ( + SELECT account_id, captain_assistant_id, MIN(id) AS keep_id, + ARRAY_AGG(id) AS duplicate_ids + FROM agent_bots + WHERE captain_assistant_id IS NOT NULL + GROUP BY account_id, captain_assistant_id + HAVING COUNT(*) > 1 +) canonical +WHERE rule.agent_bot_id = ANY(canonical.duplicate_ids) + AND rule.agent_bot_id <> canonical.keep_id; + +UPDATE bot_trigger_configs config +SET agent_bot_id = canonical.keep_id +FROM ( + SELECT account_id, captain_assistant_id, MIN(id) AS keep_id, + ARRAY_AGG(id) AS duplicate_ids + FROM agent_bots + WHERE captain_assistant_id IS NOT NULL + GROUP BY account_id, captain_assistant_id + HAVING COUNT(*) > 1 +) canonical +WHERE config.agent_bot_id = ANY(canonical.duplicate_ids) + AND config.agent_bot_id <> canonical.keep_id; + +DELETE FROM agent_bot_inboxes duplicate +USING agent_bot_inboxes canonical +WHERE duplicate.id > canonical.id + AND duplicate.agent_bot_id = canonical.agent_bot_id + AND duplicate.inbox_id = canonical.inbox_id; + +DELETE FROM agent_bots duplicate +USING agent_bots canonical +WHERE duplicate.id > canonical.id + AND duplicate.account_id IS NOT DISTINCT FROM canonical.account_id + AND duplicate.captain_assistant_id = canonical.captain_assistant_id + AND duplicate.captain_assistant_id IS NOT NULL; + +DELETE FROM captain_inboxes duplicate +USING captain_inboxes canonical +WHERE duplicate.id > canonical.id + AND duplicate.inbox_id = canonical.inbox_id + AND duplicate.deleted_at IS NULL + AND canonical.deleted_at IS NULL; + +CREATE UNIQUE INDEX IF NOT EXISTS idx_agent_bots_captain_assistant + ON agent_bots(account_id, captain_assistant_id); +CREATE UNIQUE INDEX IF NOT EXISTS idx_agent_bot_inboxes_bot_inbox + ON agent_bot_inboxes(agent_bot_id, inbox_id); +CREATE UNIQUE INDEX IF NOT EXISTS idx_captain_inboxes_active_inbox + ON captain_inboxes(inbox_id) WHERE deleted_at IS NULL;