package autoassignment import ( "context" "encoding/json" "errors" "fmt" "sync" "sync/atomic" "testing" "time" "github.com/gochat/gochat/internal/channel" "github.com/gochat/gochat/internal/model" "github.com/gochat/gochat/internal/worker" "github.com/stretchr/testify/require" "gorm.io/gorm" ) func TestAssignmentServiceAssignsOnlyTeamMember(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, nonMember, inbox, conversation := seedAssignableConversation_Cov9(t, db) teamMember := &model.User{AccountID: account.ID, Name: "team agent", Email: "team-agent@example.com", Password: "p", Active: true, Available: true} require.NoError(t, db.Create(teamMember).Error) require.NoError(t, db.Create(&model.AccountUser{AccountID: account.ID, UserID: teamMember.ID, Role: "agent"}).Error) require.NoError(t, db.Create(&model.InboxMember{InboxID: inbox.ID, UserID: teamMember.ID, Role: "agent"}).Error) team := &model.Team{AccountID: account.ID, Name: "support", AllowAutoAssignment: true} require.NoError(t, db.Create(team).Error) require.NoError(t, db.Create(&model.TeamMember{TeamID: team.ID, UserID: teamMember.ID}).Error) require.NoError(t, db.Model(conversation).Update("team_id", team.ID).Error) assigned, err := NewAssignmentService(db, rdb).AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) require.NoError(t, err) require.Equal(t, teamMember.ID, assigned) require.NotEqual(t, nonMember.ID, assigned) } func TestAssignmentServiceSkipsTeamWithoutAutoAssignment(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, agent, inbox, conversation := seedAssignableConversation_Cov9(t, db) team := &model.Team{AccountID: account.ID, Name: "manual only", AllowAutoAssignment: true} require.NoError(t, db.Create(team).Error) require.NoError(t, db.Model(team).Update("allow_auto_assignment", false).Error) require.NoError(t, db.Create(&model.TeamMember{TeamID: team.ID, UserID: agent.ID}).Error) require.NoError(t, db.Model(conversation).Update("team_id", team.ID).Error) assigned, err := NewAssignmentService(db, rdb).AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) require.NoError(t, err) require.Zero(t, assigned) } func TestAssignmentServiceBulkSkipsSoftDeletedMemberships(t *testing.T) { t.Run("inbox member", func(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, agent, inbox, conversation := seedAssignableConversation_Cov9(t, db) require.NoError(t, db.Where("inbox_id = ? AND user_id = ?", inbox.ID, agent.ID).Delete(&model.InboxMember{}).Error) assigned, err := NewAssignmentService(db, rdb).AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.Empty(t, assigned) require.NoError(t, db.First(conversation, conversation.ID).Error) require.Nil(t, conversation.AssigneeID) }) t.Run("team member and team", func(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, agent, inbox, conversation := seedAssignableConversation_Cov9(t, db) team := &model.Team{AccountID: account.ID, Name: "removed", AllowAutoAssignment: true} require.NoError(t, db.Create(team).Error) member := &model.TeamMember{TeamID: team.ID, UserID: agent.ID} require.NoError(t, db.Create(member).Error) require.NoError(t, db.Model(conversation).Update("team_id", team.ID).Error) require.NoError(t, db.Delete(member).Error) service := NewAssignmentService(db, rdb) assigned, err := service.AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.Empty(t, assigned) require.NoError(t, db.Create(&model.TeamMember{TeamID: team.ID, UserID: agent.ID}).Error) require.NoError(t, db.Delete(team).Error) assigned, err = service.AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.Empty(t, assigned) require.NoError(t, db.First(conversation, conversation.ID).Error) require.Nil(t, conversation.AssigneeID) }) } func TestAssignmentJobDrainsBacklogPastBatchLimit(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, _, inbox, first := seedAssignableConversation_Cov9(t, db) conversations := make([]model.Conversation, assignmentBatchLimit) for i := range conversations { conversations[i] = model.Conversation{ AccountID: account.ID, InboxID: inbox.ID, ContactID: first.ContactID, Status: string(model.ConversationStatusOpen), } } require.NoError(t, db.Create(&conversations).Error) policy := &model.AssignmentPolicy{AccountID: account.ID, Name: "batch", FairDistributionLimit: assignmentBatchLimit + 1, Enabled: true} require.NoError(t, db.Create(policy).Error) require.NoError(t, db.Create(&model.InboxAssignmentPolicy{InboxID: inbox.ID, AssignmentPolicyID: policy.ID}).Error) workerPool := worker.NewWorkerPool(db) job := NewAssignmentJob(NewAssignmentService(db, rdb), workerPool, rdb) enqueued, err := job.EnqueueForInbox(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.True(t, enqueued) processed, err := workerPool.ProcessOne(context.Background()) require.NoError(t, err) require.True(t, processed) require.EqualValues(t, 1, rdb.Exists(context.Background(), job.lockKey(inbox.ID)).Val(), "successor owns the original token") processed, err = workerPool.ProcessOne(context.Background()) require.NoError(t, err) require.True(t, processed) var remaining int64 require.NoError(t, db.Model(&model.Conversation{}). Where("account_id = ? AND inbox_id = ? AND assignee_id IS NULL", account.ID, inbox.ID). Count(&remaining).Error) require.Zero(t, remaining) require.EqualValues(t, 0, rdb.Exists(context.Background(), job.lockKey(inbox.ID)).Val()) var jobCount int64 require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeAssignmentJob).Count(&jobCount).Error) require.EqualValues(t, 2, jobCount) } func TestAssignmentListenerEnqueuesDurableJob(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, agent, inbox, conversation := seedAssignableConversation_Cov9(t, db) workerPool := worker.NewWorkerPool(db) listener := NewAutoAssignmentListener(db, rdb, workerPool) event := channel.NewChannelEvent(channel.EventConversationCreated, channel.ChannelWebWidget, account.ID, inbox.ID) event.ConversationID = conversation.ID require.NoError(t, listener.OnEvent(context.Background(), event)) require.NoError(t, db.First(conversation, conversation.ID).Error) require.Nil(t, conversation.AssigneeID, "event request must not scan the backlog synchronously") var backgroundJob model.BackgroundJob require.NoError(t, db.Where("job_type = ?", TaskTypeAssignmentJob).First(&backgroundJob).Error) require.Equal(t, model.BackgroundJobStatusQueued, backgroundJob.Status) var payload assignmentJobPayload require.NoError(t, json.Unmarshal(backgroundJob.Payload, &payload)) require.Equal(t, inbox.ID, payload.InboxID) require.Equal(t, account.ID, payload.AccountID) require.NotEmpty(t, payload.Token) require.NoError(t, listener.job.perform(context.Background(), &backgroundJob)) require.NoError(t, db.First(conversation, conversation.ID).Error) require.Equal(t, &agent.ID, conversation.AssigneeID) require.EqualValues(t, 0, rdb.Exists(context.Background(), listener.job.lockKey(inbox.ID)).Val()) } func TestAssignmentJobCoalescesConcurrentInboxEnqueues(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) workerPool := worker.NewWorkerPool(db) job := NewAssignmentJob(NewAssignmentService(db, rdb), workerPool, rdb) const callers = 32 start := make(chan struct{}) errs := make(chan error, callers) var created atomic.Int32 var wg sync.WaitGroup for i := 0; i < callers; i++ { wg.Add(1) go func() { defer wg.Done() <-start enqueued, err := job.EnqueueForInbox(context.Background(), 42, 7) if err != nil { errs <- err return } if enqueued { created.Add(1) } }() } close(start) wg.Wait() close(errs) for err := range errs { require.NoError(t, err) } require.EqualValues(t, 1, created.Load()) var count int64 require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeAssignmentJob).Count(&count).Error) require.EqualValues(t, 1, count) require.NoError(t, job.release(context.Background(), 42, fmt.Sprintf("not-%s", rdb.Get(context.Background(), job.lockKey(42)).Val()))) require.EqualValues(t, 1, rdb.Exists(context.Background(), job.lockKey(42)).Val(), "a stale job must not release a newer claim") } func TestAssignmentJobKeepsTokenUntilRetryDeadLetters(t *testing.T) { db, _ := setupFullAADB_Cov9(t) redisServer, rdb := setupTestRedis(t) account, _, inbox, _ := seedAssignableConversation_Cov9(t, db) now := time.Now() workerPool := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }), worker.WithBackoff(func(int) time.Duration { return time.Hour }), ) job := NewAssignmentJob(NewAssignmentService(db, rdb), workerPool, rdb) transient := errors.New("temporary database failure") require.NoError(t, db.Callback().Query().Before("gorm:query").Register("test:auto_assignment_transient", func(tx *gorm.DB) { if tx.Statement.Table == "inboxes" { tx.AddError(transient) } })) t.Cleanup(func() { _ = db.Callback().Query().Remove("test:auto_assignment_transient") }) enqueued, err := job.EnqueueForInbox(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.True(t, enqueued) processed, err := workerPool.ProcessOne(context.Background()) require.True(t, processed) require.ErrorIs(t, err, transient) token := rdb.Get(context.Background(), job.lockKey(inbox.ID)).Val() require.NotEmpty(t, token) require.Greater(t, redisServer.TTL(job.lockKey(inbox.ID)), time.Hour) redisServer.FastForward(assignmentInFlightTTL + time.Minute) enqueued, err = job.EnqueueForInbox(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.False(t, enqueued, "an event during retry backoff must coalesce") var jobCount int64 require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeAssignmentJob).Count(&jobCount).Error) require.EqualValues(t, 1, jobCount) require.Equal(t, token, rdb.Get(context.Background(), job.lockKey(inbox.ID)).Val()) for range 2 { now = now.Add(time.Hour) processed, err = workerPool.ProcessOne(context.Background()) require.True(t, processed) require.ErrorIs(t, err, transient) } var backgroundJob model.BackgroundJob require.NoError(t, db.Where("job_type = ?", TaskTypeAssignmentJob).First(&backgroundJob).Error) require.Equal(t, model.BackgroundJobStatusDead, backgroundJob.Status) require.EqualValues(t, 0, rdb.Exists(context.Background(), job.lockKey(inbox.ID)).Val()) } func TestAssignmentJobReleasesTokenOnlyAfterDeadIsPersisted(t *testing.T) { db, rdb := setupFullAADB_Cov9(t) account, _, inbox, _ := seedAssignableConversation_Cov9(t, db) workerPool := worker.NewWorkerPool(db) job := NewAssignmentJob(NewAssignmentService(db, rdb), workerPool, rdb) transient := errors.New("temporary database failure") require.NoError(t, db.Callback().Query().Before("gorm:query").Register("test:auto_assignment_terminal_failure", func(tx *gorm.DB) { if tx.Statement.Table == "inboxes" { tx.AddError(transient) } })) t.Cleanup(func() { _ = db.Callback().Query().Remove("test:auto_assignment_terminal_failure") }) enqueued, err := job.EnqueueForInbox(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.True(t, enqueued) require.NoError(t, db.Model(&model.BackgroundJob{}). Where("job_type = ?", TaskTypeAssignmentJob). Update("attempts", 2).Error) deadUpdateStarted := make(chan struct{}) allowDeadUpdate := make(chan struct{}) var once sync.Once require.NoError(t, db.Callback().Update().Before("gorm:update").Register("test:pause_assignment_dead", func(tx *gorm.DB) { updates, ok := tx.Statement.Dest.(map[string]any) if !ok || updates["status"] != model.BackgroundJobStatusDead { return } once.Do(func() { close(deadUpdateStarted) }) <-allowDeadUpdate })) t.Cleanup(func() { _ = db.Callback().Update().Remove("test:pause_assignment_dead") }) result := make(chan error, 1) go func() { processed, processErr := workerPool.ProcessOne(context.Background()) if !processed { result <- errors.New("expected terminal job to be processed") return } result <- processErr }() <-deadUpdateStarted enqueued, err = job.EnqueueForInbox(context.Background(), inbox.ID, account.ID) require.NoError(t, err) require.False(t, enqueued, "an event before dead persistence must coalesce") require.EqualValues(t, 1, rdb.Exists(context.Background(), job.lockKey(inbox.ID)).Val()) close(allowDeadUpdate) require.ErrorIs(t, <-result, transient) var backgroundJob model.BackgroundJob require.NoError(t, db.Where("job_type = ?", TaskTypeAssignmentJob).First(&backgroundJob).Error) require.Equal(t, model.BackgroundJobStatusDead, backgroundJob.Status) require.EqualValues(t, 0, rdb.Exists(context.Background(), job.lockKey(inbox.ID)).Val()) }