203 lines
7.9 KiB
Go
203 lines
7.9 KiB
Go
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
|
|
}
|