feat(conversations): queue maintenance jobs

This commit is contained in:
2026-06-05 17:25:22 +08:00
parent 6851e20aa7
commit 9a1813485f
6 changed files with 472 additions and 23 deletions
+4
View File
@@ -106,6 +106,10 @@ func Bootstrap(env string) (*App, error) {
workerPool := worker.NewWorkerPool(db)
automation.RegisterActionDeliveryJobs(workerPool, &dbProvider{db: db})
service.RegisterConversationMaintenanceJobs(workerPool, db)
if _, err := service.EnqueueScheduledItemsTrigger(context.Background(), workerPool, time.Now()); err != nil {
applogger.L().Warnf("failed to enqueue initial scheduled item trigger: %v", err)
}
// Step 5: Connect to Redis (ref: Chatwoot config/cable.yml)
rdb, err := NewRedisClient(&cfg.Redis)
+19 -18
View File
@@ -10,8 +10,9 @@ import (
type CampaignStatus string
const (
CampaignStatusActive CampaignStatus = "active"
CampaignStatusCompleted CampaignStatus = "completed"
CampaignStatusActive CampaignStatus = "active"
CampaignStatusCompleted CampaignStatus = "completed"
CampaignStatusProcessing CampaignStatus = "processing"
)
// CampaignType represents the type of campaign.
@@ -29,21 +30,21 @@ const (
// trigger_only_during_business_hours.
type Campaign struct {
model.Base
AccountID uint `gorm:"index;not null" json:"account_id"`
InboxID uint `gorm:"index;not null" json:"inbox_id"`
SenderID *uint `gorm:"index" json:"sender_id,omitempty"`
DisplayID uint `gorm:"uniqueIndex:idx_campaign_display;not null" json:"display_id"`
Title string `gorm:"size:255;not null" json:"title"`
Message string `gorm:"type:text;not null" json:"message"`
Description string `gorm:"type:text" json:"description"`
CampaignStatus CampaignStatus `gorm:"size:50;index;default:active" json:"campaign_status"`
CampaignType CampaignType `gorm:"size:50;not null" json:"campaign_type"`
Audience string `gorm:"type:jsonb;default:'{}'" json:"audience"`
TriggerRules string `gorm:"type:jsonb;default:'{}'" json:"trigger_rules"`
TemplateParams string `gorm:"type:jsonb;default:'{}'" json:"template_params"`
ScheduledAt *time.Time `gorm:"index" json:"scheduled_at,omitempty"`
Enabled bool `gorm:"default:true" json:"enabled"`
TriggerOnlyDuringBusinessHours bool `gorm:"default:false" json:"trigger_only_during_business_hours"`
AccountID uint `gorm:"index;not null" json:"account_id"`
InboxID uint `gorm:"index;not null" json:"inbox_id"`
SenderID *uint `gorm:"index" json:"sender_id,omitempty"`
DisplayID uint `gorm:"uniqueIndex:idx_campaign_display;not null" json:"display_id"`
Title string `gorm:"size:255;not null" json:"title"`
Message string `gorm:"type:text;not null" json:"message"`
Description string `gorm:"type:text" json:"description"`
CampaignStatus CampaignStatus `gorm:"size:50;index;default:active" json:"campaign_status"`
CampaignType CampaignType `gorm:"size:50;not null" json:"campaign_type"`
Audience string `gorm:"type:jsonb;default:'{}'" json:"audience"`
TriggerRules string `gorm:"type:jsonb;default:'{}'" json:"trigger_rules"`
TemplateParams string `gorm:"type:jsonb;default:'{}'" json:"template_params"`
ScheduledAt *time.Time `gorm:"index" json:"scheduled_at,omitempty"`
Enabled bool `gorm:"default:true" json:"enabled"`
TriggerOnlyDuringBusinessHours bool `gorm:"default:false" json:"trigger_only_during_business_hours"`
}
func (Campaign) TableName() string { return "campaigns" }
func (Campaign) TableName() string { return "campaigns" }
+2 -1
View File
@@ -115,6 +115,7 @@ func (b *CampaignConversationBuilder) Build(ctx context.Context, campaign *Campa
AccountID: campaign.AccountID,
InboxID: campaign.InboxID,
ContactID: contactID,
CampaignID: &campaign.ID,
Status: "open",
ChannelType: "campaign",
}
@@ -154,4 +155,4 @@ func (b *CampaignConversationBuilder) Build(ctx context.Context, campaign *Campa
}
return nil
}
}
@@ -0,0 +1,202 @@
package service
import (
"context"
"encoding/json"
"fmt"
"sync"
"time"
"github.com/gochat/gochat/internal/campaign"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
"gorm.io/gorm"
)
const (
TaskTypeScheduledTriggerItems = "scheduled:trigger_items"
TaskTypeCampaignTriggerOneoff = "campaign:trigger_oneoff"
TaskTypeConversationReopenSnoozed = "conversation:reopen_snoozed"
TaskTypeConversationResolutionScheduler = "account:conversations_resolution_scheduler"
TaskTypeConversationResolutionForAccount = "conversation:resolution"
)
const (
scheduledItemsInterval = time.Hour
scheduledItemsLookback = 3 * 24 * time.Hour
conversationResolutionLimit = 100
)
type campaignTriggerOneoffJob struct {
CampaignID uint `json:"campaign_id"`
}
type conversationResolutionJob struct {
AccountID uint `json:"account_id"`
}
var conversationMaintenanceRegistrations sync.Map
// RegisterConversationMaintenanceJobs wires Chatwoot scheduled maintenance jobs
// into the durable worker: TriggerScheduledItemsJob fans out to campaign,
// snooze-reopen, and auto-resolution jobs.
func RegisterConversationMaintenanceJobs(wp *worker.WorkerPool, db *gorm.DB) {
registerConversationMaintenanceJobsWithNow(wp, db, time.Now)
}
func registerConversationMaintenanceJobsWithNow(wp *worker.WorkerPool, db *gorm.DB, now func() time.Time) {
if wp == nil || db == nil {
return
}
if _, loaded := conversationMaintenanceRegistrations.LoadOrStore(wp, struct{}{}); loaded {
return
}
runner := &conversationMaintenanceRunner{wp: wp, db: db, now: now}
wp.Register(TaskTypeScheduledTriggerItems, runner.performScheduledTriggerItems)
wp.Register(TaskTypeCampaignTriggerOneoff, runner.performCampaignTriggerOneoff)
wp.Register(TaskTypeConversationReopenSnoozed, runner.performReopenSnoozed)
wp.Register(TaskTypeConversationResolutionScheduler, runner.performResolutionScheduler)
wp.Register(TaskTypeConversationResolutionForAccount, runner.performResolutionForAccount)
}
func EnqueueScheduledItemsTrigger(ctx context.Context, wp *worker.WorkerPool, scheduledAt time.Time) (*model.BackgroundJob, error) {
if wp == nil {
return nil, nil
}
return wp.Enqueue(ctx, TaskTypeScheduledTriggerItems, nil,
worker.WithQueue("scheduled_jobs"),
worker.WithScheduledAt(scheduledAt),
worker.WithMaxAttempts(3),
worker.WithIdempotencyKey(scheduledItemsIdempotencyKey(scheduledAt)),
)
}
func scheduledItemsIdempotencyKey(scheduledAt time.Time) string {
bucket := scheduledAt.UTC().Truncate(scheduledItemsInterval).Unix()
return fmt.Sprintf("scheduled:trigger_items:%d", bucket)
}
type conversationMaintenanceRunner struct {
wp *worker.WorkerPool
db *gorm.DB
now func() time.Time
}
func (r *conversationMaintenanceRunner) performScheduledTriggerItems(ctx context.Context, job *model.BackgroundJob) error {
now := r.now()
var campaignIDs []uint
if err := r.db.WithContext(ctx).Model(&campaign.Campaign{}).
Where("campaign_type = ? AND campaign_status = ? AND enabled = ?", campaign.CampaignTypeOneOff, campaign.CampaignStatusActive, true).
Where("scheduled_at BETWEEN ? AND ?", now.Add(-scheduledItemsLookback), now).
Pluck("id", &campaignIDs).Error; err != nil {
return fmt.Errorf("find due one-off campaigns: %w", err)
}
for _, campaignID := range campaignIDs {
_, err := r.wp.Enqueue(ctx, TaskTypeCampaignTriggerOneoff, campaignTriggerOneoffJob{CampaignID: campaignID},
worker.WithQueue("low"),
worker.WithMaxAttempts(3),
worker.WithIdempotencyKey(fmt.Sprintf("campaign:trigger_oneoff:%d", campaignID)),
)
if err != nil {
return fmt.Errorf("enqueue campaign %d: %w", campaignID, err)
}
}
if _, err := r.wp.Enqueue(ctx, TaskTypeConversationReopenSnoozed, nil, worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
return fmt.Errorf("enqueue reopen snoozed conversations: %w", err)
}
if _, err := r.wp.Enqueue(ctx, TaskTypeConversationResolutionScheduler, nil, worker.WithQueue("scheduled_jobs"), worker.WithMaxAttempts(3)); err != nil {
return fmt.Errorf("enqueue conversation resolution scheduler: %w", err)
}
_, err := EnqueueScheduledItemsTrigger(ctx, r.wp, now.Add(scheduledItemsInterval))
return err
}
func (r *conversationMaintenanceRunner) performCampaignTriggerOneoff(ctx context.Context, job *model.BackgroundJob) error {
var payload campaignTriggerOneoffJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal one-off campaign job: %w", err)
}
if payload.CampaignID == 0 {
return fmt.Errorf("invalid one-off campaign job payload: %#v", payload)
}
claimed := false
if err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
result := tx.Model(&campaign.Campaign{}).
Where("id = ? AND campaign_type = ? AND campaign_status = ? AND enabled = ?", payload.CampaignID, campaign.CampaignTypeOneOff, campaign.CampaignStatusActive, true).
Update("campaign_status", campaign.CampaignStatusProcessing)
if result.Error != nil {
return result.Error
}
claimed = result.RowsAffected == 1
return nil
}); err != nil {
return fmt.Errorf("claim one-off campaign %d: %w", payload.CampaignID, err)
}
if !claimed {
return nil
}
if err := campaign.NewCampaignService(r.db).TriggerCampaign(ctx, payload.CampaignID); err != nil {
_ = r.db.WithContext(ctx).Model(&campaign.Campaign{}).Where("id = ?", payload.CampaignID).Update("campaign_status", campaign.CampaignStatusActive).Error
return err
}
return r.db.WithContext(ctx).Model(&campaign.Campaign{}).Where("id = ?", payload.CampaignID).Update("campaign_status", campaign.CampaignStatusCompleted).Error
}
func (r *conversationMaintenanceRunner) performReopenSnoozed(ctx context.Context, job *model.BackgroundJob) error {
now := r.now()
nowUnix := now.Unix()
lookbackUnix := now.Add(-scheduledItemsLookback).Unix()
updates := map[string]any{
"status": string(model.ConversationStatusOpen),
"snoozed_until": nil,
"resumed_at": now,
}
return r.db.WithContext(ctx).Model(&model.Conversation{}).
Where("status = ?", string(model.ConversationStatusSnoozed)).
Where("snoozed_until BETWEEN ? AND ?", lookbackUnix, nowUnix).
Updates(updates).Error
}
func (r *conversationMaintenanceRunner) performResolutionScheduler(ctx context.Context, job *model.BackgroundJob) error {
var accountIDs []uint
if err := r.db.WithContext(ctx).Model(&model.Account{}).
Where("auto_resolve_duration > 0").
Pluck("id", &accountIDs).Error; err != nil {
return fmt.Errorf("find auto-resolve accounts: %w", err)
}
for _, accountID := range accountIDs {
if _, err := r.wp.Enqueue(ctx, TaskTypeConversationResolutionForAccount, conversationResolutionJob{AccountID: accountID}, worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
return fmt.Errorf("enqueue account conversation resolution %d: %w", accountID, err)
}
}
return nil
}
func (r *conversationMaintenanceRunner) performResolutionForAccount(ctx context.Context, job *model.BackgroundJob) error {
var payload conversationResolutionJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal conversation resolution job: %w", err)
}
if payload.AccountID == 0 {
return fmt.Errorf("invalid conversation resolution job payload: %#v", payload)
}
var account model.Account
if err := r.db.WithContext(ctx).First(&account, payload.AccountID).Error; err != nil {
return fmt.Errorf("load auto-resolve account %d: %w", payload.AccountID, err)
}
if account.AutoResolveDuration <= 0 {
return nil
}
cutoff := r.now().Add(-time.Duration(account.AutoResolveDuration) * time.Minute).Unix()
now := r.now()
updates := map[string]any{
"status": string(model.ConversationStatusResolved),
"resolved_at": now,
}
return r.db.WithContext(ctx).Model(&model.Conversation{}).
Where("account_id = ? AND status = ? AND contact_id <> 0", account.ID, string(model.ConversationStatusOpen)).
Where("last_activity_at IS NOT NULL AND last_activity_at < ?", cutoff).
Limit(conversationResolutionLimit).
Updates(updates).Error
}
@@ -0,0 +1,221 @@
package service
import (
"context"
"fmt"
"testing"
"time"
"github.com/gochat/gochat/internal/campaign"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
"gorm.io/gorm"
)
func TestConversationMaintenanceJobsTriggerScheduledItemsFanOut(t *testing.T) {
now := time.Date(2026, 6, 5, 19, 0, 0, 0, time.UTC)
db := setupServiceTestDB(t)
if err := db.AutoMigrate(&campaign.Campaign{}); err != nil {
t.Fatalf("migrate campaign: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }))
registerConversationMaintenanceJobsWithNow(wp, db, func() time.Time { return now })
account := createTestAccount(t, db)
inbox := createTestInbox(t, db, account.ID, "sms")
contact := createTestContact(t, db, account.ID)
dueAt := now.Add(-time.Hour)
futureAt := now.Add(time.Hour)
dueCampaign := createTestOneoffCampaign(t, db, account.ID, inbox.ID, contact.ID, dueAt)
createTestOneoffCampaign(t, db, account.ID, inbox.ID, contact.ID, futureAt)
if _, err := EnqueueScheduledItemsTrigger(context.Background(), wp, now); err != nil {
t.Fatalf("enqueue scheduled trigger: %v", err)
}
processed, err := wp.ProcessOne(context.Background())
if err != nil || !processed {
t.Fatalf("process scheduled trigger: processed=%v err=%v", processed, err)
}
assertJobCount(t, db, TaskTypeCampaignTriggerOneoff, 1)
assertJobCount(t, db, TaskTypeConversationReopenSnoozed, 1)
assertJobCount(t, db, TaskTypeConversationResolutionScheduler, 1)
var campaignJob model.BackgroundJob
if err := db.Where("job_type = ?", TaskTypeCampaignTriggerOneoff).First(&campaignJob).Error; err != nil {
t.Fatalf("load campaign job: %v", err)
}
if want := fmt.Sprintf("campaign:trigger_oneoff:%d", dueCampaign.ID); campaignJob.IdempotencyKey != want {
t.Fatalf("expected due campaign idempotency key %q, got %q", want, campaignJob.IdempotencyKey)
}
var nextTrigger model.BackgroundJob
if err := db.Where("job_type = ? AND status = ?", TaskTypeScheduledTriggerItems, model.BackgroundJobStatusQueued).First(&nextTrigger).Error; err != nil {
t.Fatalf("load next trigger: %v", err)
}
if !nextTrigger.ScheduledAt.Equal(now.Add(scheduledItemsInterval)) {
t.Fatalf("expected next trigger at %s, got %s", now.Add(scheduledItemsInterval), nextTrigger.ScheduledAt)
}
}
func TestConversationMaintenanceJobsProcessCampaignSnoozeAndResolution(t *testing.T) {
now := time.Date(2026, 6, 5, 20, 0, 0, 0, time.UTC)
db := setupServiceTestDB(t)
if err := db.AutoMigrate(&campaign.Campaign{}); err != nil {
t.Fatalf("migrate campaign: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }))
registerConversationMaintenanceJobsWithNow(wp, db, func() time.Time { return now })
account := createTestAccount(t, db)
account.AutoResolveDuration = 30
if err := db.Save(account).Error; err != nil {
t.Fatalf("save auto resolve account: %v", err)
}
inbox := createTestInbox(t, db, account.ID, "sms")
contact := createTestContact(t, db, account.ID)
dueCampaign := createTestOneoffCampaign(t, db, account.ID, inbox.ID, contact.ID, now.Add(-time.Hour))
dueSnooze := now.Add(-time.Minute).Unix()
snoozed := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
if err := db.Model(snoozed).Updates(map[string]any{"status": string(model.ConversationStatusSnoozed), "snoozed_until": dueSnooze}).Error; err != nil {
t.Fatalf("snooze conversation: %v", err)
}
oldActivity := now.Add(-45 * time.Minute).Unix()
oldOpen := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
if err := db.Model(oldOpen).Update("last_activity_at", oldActivity).Error; err != nil {
t.Fatalf("set old activity: %v", err)
}
recentActivity := now.Add(-5 * time.Minute).Unix()
recentOpen := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
if err := db.Model(recentOpen).Update("last_activity_at", recentActivity).Error; err != nil {
t.Fatalf("set recent activity: %v", err)
}
if _, err := wp.Enqueue(context.Background(), TaskTypeCampaignTriggerOneoff, campaignTriggerOneoffJob{CampaignID: dueCampaign.ID}, worker.WithQueue("low")); err != nil {
t.Fatalf("enqueue campaign: %v", err)
}
if _, err := wp.Enqueue(context.Background(), TaskTypeConversationReopenSnoozed, nil, worker.WithQueue("low")); err != nil {
t.Fatalf("enqueue reopen: %v", err)
}
if _, err := wp.Enqueue(context.Background(), TaskTypeConversationResolutionScheduler, nil, worker.WithQueue("scheduled_jobs")); err != nil {
t.Fatalf("enqueue scheduler: %v", err)
}
processRequiredJob(t, wp, "campaign")
processRequiredJob(t, wp, "reopen")
processRequiredJob(t, wp, "scheduler")
processRequiredJob(t, wp, "resolution")
var completed campaign.Campaign
if err := db.First(&completed, dueCampaign.ID).Error; err != nil {
t.Fatalf("load campaign: %v", err)
}
if completed.CampaignStatus != campaign.CampaignStatusCompleted {
t.Fatalf("expected completed campaign, got %s", completed.CampaignStatus)
}
var campaignMessages int64
if err := db.Model(&model.Message{}).Where("content = ?", dueCampaign.Message).Count(&campaignMessages).Error; err != nil {
t.Fatalf("count campaign messages: %v", err)
}
if campaignMessages != 1 {
t.Fatalf("expected one campaign message, got %d", campaignMessages)
}
if _, err := wp.Enqueue(context.Background(), TaskTypeCampaignTriggerOneoff, campaignTriggerOneoffJob{CampaignID: dueCampaign.ID}, worker.WithQueue("low")); err != nil {
t.Fatalf("enqueue duplicate campaign: %v", err)
}
processRequiredJob(t, wp, "duplicate campaign")
if err := db.Model(&model.Message{}).Where("content = ?", dueCampaign.Message).Count(&campaignMessages).Error; err != nil {
t.Fatalf("count duplicate campaign messages: %v", err)
}
if campaignMessages != 1 {
t.Fatalf("expected duplicate campaign job to be idempotent, got %d messages", campaignMessages)
}
var reopened model.Conversation
if err := db.First(&reopened, snoozed.ID).Error; err != nil {
t.Fatalf("load reopened conversation: %v", err)
}
if reopened.Status != string(model.ConversationStatusOpen) || reopened.SnoozedUntil != nil || reopened.ResumedAt == nil {
t.Fatalf("expected snoozed conversation reopened, got status=%s snoozed=%v resumed=%v", reopened.Status, reopened.SnoozedUntil, reopened.ResumedAt)
}
var resolved model.Conversation
if err := db.First(&resolved, oldOpen.ID).Error; err != nil {
t.Fatalf("load resolved conversation: %v", err)
}
if resolved.Status != string(model.ConversationStatusResolved) || resolved.ResolvedAt == nil {
t.Fatalf("expected old open conversation resolved, got status=%s resolved_at=%v", resolved.Status, resolved.ResolvedAt)
}
var recent model.Conversation
if err := db.First(&recent, recentOpen.ID).Error; err != nil {
t.Fatalf("load recent conversation: %v", err)
}
if recent.Status != string(model.ConversationStatusOpen) {
t.Fatalf("expected recent conversation to remain open, got %s", recent.Status)
}
}
func TestConversationMaintenanceJobsRetryMissingResolutionAccount(t *testing.T) {
now := time.Date(2026, 6, 5, 21, 0, 0, 0, time.UTC)
db := setupServiceTestDB(t)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }), worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }))
registerConversationMaintenanceJobsWithNow(wp, db, func() time.Time { return now })
if _, err := wp.Enqueue(context.Background(), TaskTypeConversationResolutionForAccount, conversationResolutionJob{AccountID: 9999}, worker.WithQueue("low"), worker.WithMaxAttempts(3)); err != nil {
t.Fatalf("enqueue missing account: %v", err)
}
processed, err := wp.ProcessOne(context.Background())
if err == nil || !processed {
t.Fatalf("expected missing account to retry, processed=%v err=%v", processed, err)
}
var job model.BackgroundJob
if err := db.Where("job_type = ?", TaskTypeConversationResolutionForAccount).First(&job).Error; err != nil {
t.Fatalf("load resolution job: %v", err)
}
if job.Status != model.BackgroundJobStatusRetrying || job.LastError == "" {
t.Fatalf("expected retrying resolution job with error, got status=%s last_error=%q", job.Status, job.LastError)
}
}
func createTestOneoffCampaign(t *testing.T, db *gorm.DB, accountID, inboxID, contactID uint, scheduledAt time.Time) *campaign.Campaign {
t.Helper()
c := &campaign.Campaign{
AccountID: accountID,
InboxID: inboxID,
DisplayID: uint(time.Now().UnixNano()),
Title: fmt.Sprintf("Campaign %d", time.Now().UnixNano()),
Message: fmt.Sprintf("Campaign message %d", time.Now().UnixNano()),
CampaignStatus: campaign.CampaignStatusActive,
CampaignType: campaign.CampaignTypeOneOff,
Audience: fmt.Sprintf(`{"contact_ids":[%d]}`, contactID),
TriggerRules: `{}`,
TemplateParams: `{}`,
ScheduledAt: &scheduledAt,
Enabled: true,
}
if err := db.Create(c).Error; err != nil {
t.Fatalf("create campaign: %v", err)
}
return c
}
func assertJobCount(t *testing.T, db *gorm.DB, jobType string, want int64) {
t.Helper()
var count int64
if err := db.Model(&model.BackgroundJob{}).Where("job_type = ? AND status = ?", jobType, model.BackgroundJobStatusQueued).Count(&count).Error; err != nil {
t.Fatalf("count jobs %s: %v", jobType, err)
}
if count != want {
t.Fatalf("expected %d queued %s jobs, got %d", want, jobType, count)
}
}
func processRequiredJob(t *testing.T, wp *worker.WorkerPool, name string) {
t.Helper()
processed, err := wp.ProcessOne(context.Background())
if err != nil || !processed {
t.Fatalf("process %s job: processed=%v err=%v", name, processed, err)
}
}