From 391fe98ad6ef5c131171d12dd2ad477321f12cbd Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 24 Aug 2026 00:16:03 +0800 Subject: [PATCH] HH-565: complete assignment parity regressions (#142) * fix(HH-565): complete assignment parity regressions * fix(HH-565): scope capacity exclusions per agent --------- Co-authored-by: Rogee --- .../autoassignment/assignment_job_test.go | 18 ++- .../internal/autoassignment/capacity_test.go | 146 +++++++++++++++++- backend/internal/autoassignment/service.go | 89 ++++++++++- .../repository/conversationassignee/gate.go | 1 + backend/internal/worker/worker_test.go | 50 ++++++ 5 files changed, 292 insertions(+), 12 deletions(-) diff --git a/backend/internal/autoassignment/assignment_job_test.go b/backend/internal/autoassignment/assignment_job_test.go index e0816966..a012cc6b 100644 --- a/backend/internal/autoassignment/assignment_job_test.go +++ b/backend/internal/autoassignment/assignment_job_test.go @@ -49,15 +49,17 @@ func TestAssignmentServiceSkipsTeamWithoutAutoAssignment(t *testing.T) { require.Zero(t, assigned) } -func TestAssignmentServiceSkipsSoftDeletedMemberships(t *testing.T) { +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).AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) + assigned, err := NewAssignmentService(db, rdb).AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) - require.Zero(t, assigned) + 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) { @@ -71,15 +73,17 @@ func TestAssignmentServiceSkipsSoftDeletedMemberships(t *testing.T) { require.NoError(t, db.Delete(member).Error) service := NewAssignmentService(db, rdb) - assigned, err := service.AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) + assigned, err := service.AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) - require.Zero(t, assigned) + 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.AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) + assigned, err = service.AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) require.NoError(t, err) - require.Zero(t, assigned) + require.Empty(t, assigned) + require.NoError(t, db.First(conversation, conversation.ID).Error) + require.Nil(t, conversation.AssigneeID) }) } diff --git a/backend/internal/autoassignment/capacity_test.go b/backend/internal/autoassignment/capacity_test.go index 81950403..ad350f6f 100644 --- a/backend/internal/autoassignment/capacity_test.go +++ b/backend/internal/autoassignment/capacity_test.go @@ -12,6 +12,7 @@ import ( "gorm.io/gorm/logger" "github.com/gochat/gochat/internal/model" + "github.com/gochat/gochat/internal/repository" ) func setupCapacityAssignmentDB(t *testing.T) *gorm.DB { @@ -30,6 +31,8 @@ func setupCapacityAssignmentDB(t *testing.T) *gorm.DB { &model.InboxAssignmentPolicy{}, &model.AgentCapacityPolicy{}, &model.InboxCapacityLimit{}, + &model.Tag{}, + &model.ConversationLabel{}, )) return db } @@ -79,13 +82,14 @@ func TestAssignmentServiceIgnoresCapacityWhenAdvancedAssignmentIsDisabled(t *tes require.NoError(t, db.Create(inbox).Error) agent := createCapacityUser(t, db, account.ID, "default@test.com") require.NoError(t, db.Create(&model.InboxMember{InboxID: inbox.ID, UserID: agent.ID}).Error) - policy := &model.AgentCapacityPolicy{AccountID: account.ID, Name: "Excluded", ExclusionRules: []byte(`{}`)} + policy := &model.AgentCapacityPolicy{AccountID: account.ID, Name: "Excluded", ExclusionRules: []byte(`{"exclude_older_than_hours":1}`)} require.NoError(t, db.Create(policy).Error) require.NoError(t, db.Model(&model.AccountUser{}).Where("account_id = ? AND user_id = ?", account.ID, agent.ID).Update("agent_capacity_policy_id", policy.ID).Error) require.NoError(t, db.Create(&model.InboxCapacityLimit{AgentCapacityPolicyID: policy.ID, InboxID: inbox.ID, ConversationLimit: 0}).Error) contact := &model.Contact{AccountID: account.ID, Name: "Contact"} require.NoError(t, db.Create(contact).Error) - conversation := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen)} + old := time.Now().Add(-2 * time.Hour).Unix() + conversation := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &old} require.NoError(t, db.Create(conversation).Error) assigned, err := NewAssignmentService(db, rdb).AssignConversation(context.Background(), conversation.ID, inbox.ID, account.ID) @@ -93,6 +97,144 @@ func TestAssignmentServiceIgnoresCapacityWhenAdvancedAssignmentIsDisabled(t *tes require.Equal(t, agent.ID, assigned) } +func TestAssignmentServiceAppliesCapacityPolicyExclusions(t *testing.T) { + _, rdb := setupTestRedis(t) + db := setupCapacityAssignmentDB(t) + account := &model.Account{Name: "Exclusions", FeatureFlags: `{"advanced_assignment":true}`} + require.NoError(t, db.Create(account).Error) + inbox := &model.Inbox{AccountID: account.ID, Name: "Inbox", ChannelType: string(model.InboxChannelTypeWebWidget), EnableAutoAssignment: true} + require.NoError(t, db.Create(inbox).Error) + agent := createCapacityUser(t, db, account.ID, "exclusions@test.com") + require.NoError(t, db.Create(&model.InboxMember{InboxID: inbox.ID, UserID: agent.ID}).Error) + capacityPolicy := &model.AgentCapacityPolicy{ + AccountID: account.ID, + Name: "Exclude stale spam", + ExclusionRules: []byte(`{"excluded_labels":["spam"],"exclude_older_than_hours":24}`), + } + require.NoError(t, db.Create(capacityPolicy).Error) + require.NoError(t, db.Model(&model.AccountUser{}).Where("account_id = ? AND user_id = ?", account.ID, agent.ID).Update("agent_capacity_policy_id", capacityPolicy.ID).Error) + require.NoError(t, db.Create(&model.InboxCapacityLimit{AgentCapacityPolicyID: capacityPolicy.ID, InboxID: inbox.ID, ConversationLimit: 10}).Error) + contact := &model.Contact{AccountID: account.ID, Name: "Contact"} + require.NoError(t, db.Create(contact).Error) + + now := time.Now().Unix() + old := time.Now().Add(-25 * time.Hour).Unix() + recent := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &now} + labeled := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &now} + stale := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &old} + unknownAge := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen)} + require.NoError(t, db.Create([]*model.Conversation{recent, labeled, stale, unknownAge}).Error) + tag := &model.Tag{AccountID: account.ID, Name: "spam"} + require.NoError(t, db.Create(tag).Error) + require.NoError(t, db.Create(&model.ConversationLabel{AccountID: account.ID, ConversationID: labeled.ID, TagID: tag.ID}).Error) + + assigned, err := NewAssignmentService(db, rdb).AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) + require.NoError(t, err) + require.Equal(t, []uint{recent.ID}, assigned) + for _, excluded := range []*model.Conversation{labeled, stale, unknownAge} { + require.NoError(t, db.First(excluded, excluded.ID).Error) + require.Nil(t, excluded.AssigneeID) + } +} + +func TestAssignmentServiceAppliesEachAgentsCapacityPolicyExclusions(t *testing.T) { + _, rdb := setupTestRedis(t) + db := setupCapacityAssignmentDB(t) + account := &model.Account{Name: "Per-agent exclusions", FeatureFlags: `{"advanced_assignment":true}`} + require.NoError(t, db.Create(account).Error) + inbox := &model.Inbox{AccountID: account.ID, Name: "Inbox", ChannelType: string(model.InboxChannelTypeWebWidget), EnableAutoAssignment: true} + require.NoError(t, db.Create(inbox).Error) + spamAgent := createCapacityUser(t, db, account.ID, "spam-policy@test.com") + vipAgent := createCapacityUser(t, db, account.ID, "vip-policy@test.com") + require.NoError(t, db.Create([]*model.InboxMember{{InboxID: inbox.ID, UserID: spamAgent.ID}, {InboxID: inbox.ID, UserID: vipAgent.ID}}).Error) + + spamPolicy := &model.AgentCapacityPolicy{AccountID: account.ID, Name: "Exclude spam", ExclusionRules: []byte(`{"excluded_labels":["spam"]}`)} + vipPolicy := &model.AgentCapacityPolicy{AccountID: account.ID, Name: "Exclude vip", ExclusionRules: []byte(`{"excluded_labels":["vip"]}`)} + require.NoError(t, db.Create(spamPolicy).Error) + require.NoError(t, db.Create(vipPolicy).Error) + require.NoError(t, db.Model(&model.AccountUser{}).Where("account_id = ? AND user_id = ?", account.ID, spamAgent.ID).Update("agent_capacity_policy_id", spamPolicy.ID).Error) + require.NoError(t, db.Model(&model.AccountUser{}).Where("account_id = ? AND user_id = ?", account.ID, vipAgent.ID).Update("agent_capacity_policy_id", vipPolicy.ID).Error) + require.NoError(t, db.Create([]*model.InboxCapacityLimit{ + {AgentCapacityPolicyID: spamPolicy.ID, InboxID: inbox.ID, ConversationLimit: 10}, + {AgentCapacityPolicyID: vipPolicy.ID, InboxID: inbox.ID, ConversationLimit: 10}, + }).Error) + + contact := &model.Contact{AccountID: account.ID, Name: "Contact"} + require.NoError(t, db.Create(contact).Error) + now := time.Now().Unix() + spam := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &now} + vip := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &now} + require.NoError(t, db.Create([]*model.Conversation{spam, vip}).Error) + spamTag := &model.Tag{AccountID: account.ID, Name: "spam"} + vipTag := &model.Tag{AccountID: account.ID, Name: "vip"} + require.NoError(t, db.Create([]*model.Tag{spamTag, vipTag}).Error) + require.NoError(t, db.Create([]*model.ConversationLabel{ + {AccountID: account.ID, ConversationID: spam.ID, TagID: spamTag.ID}, + {AccountID: account.ID, ConversationID: vip.ID, TagID: vipTag.ID}, + }).Error) + + assigned, err := NewAssignmentService(db, rdb).AssignUnassignedConversations(context.Background(), inbox.ID, account.ID) + require.NoError(t, err) + require.Equal(t, []uint{spam.ID, vip.ID}, assigned) + require.NoError(t, db.First(spam, spam.ID).Error) + require.NoError(t, db.First(vip, vip.ID).Error) + require.Equal(t, vipAgent.ID, *spam.AssigneeID) + require.Equal(t, spamAgent.ID, *vip.AssigneeID) +} + +func TestAssignmentEligibilityRejectsSoftDeletedUserAndAccountUser(t *testing.T) { + for _, deleted := range []string{"user", "account user"} { + t.Run(deleted, func(t *testing.T) { + db := setupCapacityAssignmentDB(t) + account := &model.Account{Name: deleted} + require.NoError(t, db.Create(account).Error) + inbox := &model.Inbox{AccountID: account.ID, Name: "Inbox", ChannelType: string(model.InboxChannelTypeWebWidget), EnableAutoAssignment: true} + require.NoError(t, db.Create(inbox).Error) + agent := createCapacityUser(t, db, account.ID, deleted+"@test.com") + require.NoError(t, db.Create(&model.InboxMember{InboxID: inbox.ID, UserID: agent.ID}).Error) + contact := &model.Contact{AccountID: account.ID, Name: "Contact"} + require.NoError(t, db.Create(contact).Error) + conversation := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen)} + require.NoError(t, db.Create(conversation).Error) + + if deleted == "user" { + require.NoError(t, db.Delete(agent).Error) + } else { + require.NoError(t, db.Where("account_id = ? AND user_id = ?", account.ID, agent.ID).Delete(&model.AccountUser{}).Error) + } + + agents, err := NewAssignmentService(db, nil).getEligibleAgents(context.Background(), inbox.ID, account.ID, nil, false) + require.NoError(t, err) + require.Empty(t, agents) + assigned, err := repository.AutoAssignConversation(context.Background(), db, account.ID, inbox.ID, conversation.ID, agent.ID) + require.NoError(t, err) + require.False(t, assigned) + }) + } +} + +func TestLongestWaitingSortsUnknownActivityLast(t *testing.T) { + db := setupCapacityAssignmentDB(t) + account := &model.Account{Name: "Longest waiting"} + require.NoError(t, db.Create(account).Error) + inbox := &model.Inbox{AccountID: account.ID, Name: "Inbox", ChannelType: string(model.InboxChannelTypeWebWidget)} + require.NoError(t, db.Create(inbox).Error) + contact := &model.Contact{AccountID: account.ID, Name: "Contact"} + require.NoError(t, db.Create(contact).Error) + olderActivity := time.Now().Add(-2 * time.Hour).Unix() + newerActivity := time.Now().Add(-time.Hour).Unix() + unknown := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen)} + newer := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &newerActivity} + older := &model.Conversation{AccountID: account.ID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), LastActivityAt: &olderActivity} + require.NoError(t, db.Create([]*model.Conversation{unknown, newer, older}).Error) + + conversations, err := NewAssignmentService(db, nil).findUnassignedConversations( + context.Background(), inbox.ID, account.ID, &model.AssignmentPolicy{ConversationPriority: 1}, + ) + require.NoError(t, err) + require.Equal(t, []uint{older.ID, newer.ID, unknown.ID}, []uint{conversations[0].ID, conversations[1].ID, conversations[2].ID}) +} + func createCapacityUser(t *testing.T, db *gorm.DB, accountID uint, email string) *model.User { t.Helper() user := &model.User{AccountID: accountID, Name: email, Email: email, Password: "hashed", Active: true, Available: true} diff --git a/backend/internal/autoassignment/service.go b/backend/internal/autoassignment/service.go index b597e4d9..30af98e2 100644 --- a/backend/internal/autoassignment/service.go +++ b/backend/internal/autoassignment/service.go @@ -109,6 +109,12 @@ func (s *AssignmentService) AssignUnassignedConversations(ctx context.Context, i if err != nil { return assignedIDs, fmt.Errorf("get eligible agents for conversation %d: %w", conv.ID, err) } + if advancedAssignment { + agents, err = s.filterAgentsByCapacityExclusions(ctx, accountID, &conv, agents) + if err != nil { + return assignedIDs, fmt.Errorf("apply capacity exclusions for conversation %d: %w", conv.ID, err) + } + } if len(agents) == 0 { continue } @@ -176,7 +182,7 @@ func (s *AssignmentService) AssignConversation(ctx context.Context, conversation var conversation model.Conversation if err := s.db.WithContext(ctx). - Select("team_id"). + Select("id", "team_id", "last_activity_at"). Where("id = ? AND inbox_id = ? AND account_id = ?", conversationID, inboxID, accountID). First(&conversation).Error; err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { @@ -190,6 +196,12 @@ func (s *AssignmentService) AssignConversation(ctx context.Context, conversation if err != nil { return 0, fmt.Errorf("get eligible agents: %w", err) } + if advancedAssignment { + agents, err = s.filterAgentsByCapacityExclusions(ctx, accountID, &conversation, agents) + if err != nil { + return 0, fmt.Errorf("apply capacity exclusions: %w", err) + } + } if len(agents) == 0 { return 0, nil } @@ -292,7 +304,10 @@ func (s *AssignmentService) findUnassignedConversations(ctx context.Context, inb query = query.Where("last_activity_at IS NULL OR last_activity_at >= ?", cutoff) } if policy != nil && policy.ConversationPriority == 1 { - query = query.Order("last_activity_at ASC").Order("created_at ASC") + // Unknown activity is not treated as longest waiting. PostgreSQL already + // puts NULL last for ASC; the CASE makes SQLite follow the same contract. + query = query.Order("CASE WHEN last_activity_at IS NULL THEN 1 ELSE 0 END ASC"). + Order("last_activity_at ASC").Order("created_at ASC") } else { query = query.Order("created_at ASC") } @@ -310,7 +325,7 @@ func (s *AssignmentService) getEligibleAgents(ctx context.Context, inboxID uint, Distinct("inbox_members.user_id"). Joins("JOIN users ON users.id = inbox_members.user_id"). Joins("JOIN account_users ON account_users.user_id = inbox_members.user_id AND account_users.account_id = ?", accountID). - Where("inbox_members.inbox_id = ? AND inbox_members.deleted_at IS NULL AND users.available = ? AND users.active = ? AND account_users.role IN ?", + Where("inbox_members.inbox_id = ? AND inbox_members.deleted_at IS NULL AND users.deleted_at IS NULL AND account_users.deleted_at IS NULL AND users.available = ? AND users.active = ? AND account_users.role IN ?", inboxID, true, true, []string{"agent", "administrator"}) if teamID != nil { query = query. @@ -341,6 +356,74 @@ func (s *AssignmentService) getEligibleAgents(ctx context.Context, inboxID uint, return availableIDs, nil } +func (s *AssignmentService) filterAgentsByCapacityExclusions(ctx context.Context, accountID uint, conversation *model.Conversation, agentIDs []uint) ([]uint, error) { + eligible := make([]uint, 0, len(agentIDs)) + for _, agentID := range agentIDs { + excluded, err := s.capacityPolicyExcludesConversation(ctx, accountID, agentID, conversation) + if err != nil { + return nil, err + } + if !excluded { + eligible = append(eligible, agentID) + } + } + return eligible, nil +} + +func (s *AssignmentService) capacityPolicyExcludesConversation(ctx context.Context, accountID, agentID uint, conversation *model.Conversation) (bool, error) { + var accountUser model.AccountUser + if err := s.db.WithContext(ctx). + Where("account_id = ? AND user_id = ?", accountID, agentID). + First(&accountUser).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return true, nil + } + return false, fmt.Errorf("load account user capacity policy: %w", err) + } + if accountUser.AgentCapacityPolicyID == nil { + return false, nil + } + + var policy model.AgentCapacityPolicy + if err := s.db.WithContext(ctx). + Where("id = ? AND account_id = ?", *accountUser.AgentCapacityPolicyID, accountID). + First(&policy).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return false, nil + } + return false, fmt.Errorf("load agent capacity policy: %w", err) + } + var rules struct { + ExcludedLabels []string `json:"excluded_labels"` + ExcludeOlderThanHours int `json:"exclude_older_than_hours"` + } + if len(policy.ExclusionRules) > 0 { + if err := json.Unmarshal(policy.ExclusionRules, &rules); err != nil { + return false, fmt.Errorf("decode capacity exclusion rules: %w", err) + } + } + if rules.ExcludeOlderThanHours > 0 { + cutoff := time.Now().Add(-time.Duration(rules.ExcludeOlderThanHours) * time.Hour).Unix() + if conversation.LastActivityAt == nil || *conversation.LastActivityAt < cutoff { + return true, nil + } + } + if len(rules.ExcludedLabels) == 0 { + return false, nil + } + + var count int64 + err := s.db.WithContext(ctx).Table("conversation_labels"). + Joins("JOIN tags ON tags.id = conversation_labels.tag_id AND tags.deleted_at IS NULL"). + Where("conversation_labels.conversation_id = ? AND conversation_labels.account_id = ? AND tags.account_id = ? AND tags.name IN ?", + conversation.ID, accountID, accountID, rules.ExcludedLabels). + Count(&count).Error + if err != nil { + return false, fmt.Errorf("check capacity exclusion labels: %w", err) + } + return count > 0, nil +} + func (s *AssignmentService) agentHasInboxCapacity(ctx context.Context, accountID, inboxID, agentID, excludeConversationID uint) (bool, error) { var accountUser model.AccountUser if err := s.db.WithContext(ctx). diff --git a/backend/internal/repository/conversationassignee/gate.go b/backend/internal/repository/conversationassignee/gate.go index ea589d18..27808113 100644 --- a/backend/internal/repository/conversationassignee/gate.go +++ b/backend/internal/repository/conversationassignee/gate.go @@ -48,5 +48,6 @@ func eligible(db *gorm.DB, assigneeID uint) *gorm.DB { Joins("JOIN users ON users.id = account_users.user_id"). Where("account_users.account_id = conversations.account_id"). Where("account_users.user_id = ?", assigneeID). + Where("account_users.deleted_at IS NULL AND users.deleted_at IS NULL"). Where("account_users.role IN ? AND users.active = ?", []string{"agent", "administrator"}, true) } diff --git a/backend/internal/worker/worker_test.go b/backend/internal/worker/worker_test.go index 7c1a045f..e21dd3f1 100644 --- a/backend/internal/worker/worker_test.go +++ b/backend/internal/worker/worker_test.go @@ -422,6 +422,56 @@ func TestRedisClaimCanceledWhenShutdownRacesClaim(t *testing.T) { require.Equal(t, model.BackgroundJobStatusQueued, loadJob(t, db, job.ID).Status) } +func TestRedisFailureLifecycleTransitionsOncePerClaim(t *testing.T) { + db := newWorkerTestDB(t) + _, rdb := newMiniRedis(t) + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + wp := newRedisWorkerPool(t, db, rdb, + WithNow(func() time.Time { return now }), + WithBackoff(func(int) time.Duration { return time.Hour }), + ) + boom := errors.New("provider timeout") + var handled, failed atomic.Int32 + wp.Register("failing_stream_job", func(context.Context, *model.BackgroundJob) error { + handled.Add(1) + return boom + }) + wp.RegisterFailureHandler("failing_stream_job", func(context.Context, *model.BackgroundJob, time.Duration) error { + failed.Add(1) + return nil + }) + wp.ensureConsumerGroups(context.Background()) + job, err := wp.Enqueue(context.Background(), "failing_stream_job", nil, WithMaxAttempts(2)) + require.NoError(t, err) + stream := wp.streamKeyFor(job.Queue) + + wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "1-0", Values: map[string]any{"job_id": job.ID}}) + retrying := loadJob(t, db, job.ID) + require.Equal(t, model.BackgroundJobStatusRetrying, retrying.Status) + require.Equal(t, 1, retrying.Attempts) + require.EqualValues(t, 1, handled.Load()) + require.EqualValues(t, 1, failed.Load()) + + wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "1-1", Values: map[string]any{"job_id": job.ID}}) + require.Equal(t, model.BackgroundJobStatusRetrying, loadJob(t, db, job.ID).Status) + require.EqualValues(t, 1, handled.Load(), "duplicate delivery during backoff must not rerun the handler") + require.EqualValues(t, 1, failed.Load(), "duplicate delivery during backoff must not repeat failure side effects") + + now = now.Add(time.Hour) + wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "2-0", Values: map[string]any{"job_id": job.ID}}) + dead := loadJob(t, db, job.ID) + require.Equal(t, model.BackgroundJobStatusDead, dead.Status) + require.Equal(t, 2, dead.Attempts) + require.NotNil(t, dead.FailedAt) + require.EqualValues(t, 2, handled.Load()) + require.EqualValues(t, 2, failed.Load()) + + wp.processRedisMessage(context.Background(), context.Background(), stream, redis.XMessage{ID: "2-1", Values: map[string]any{"job_id": job.ID}}) + require.Equal(t, model.BackgroundJobStatusDead, loadJob(t, db, job.ID).Status) + require.EqualValues(t, 2, handled.Load(), "dead jobs must not be executed again") + require.EqualValues(t, 2, failed.Load(), "dead jobs must not repeat failure side effects") +} + // TestRedisEnqueueAndProcessEndToEnd verifies the full Redis path: Enqueue // XADDs to the stream, XREADGROUP picks it up, DB claim succeeds, handler // runs, and the job reaches "completed" status.