feat(crm): queue contact bulk actions
This commit is contained in:
@@ -2,6 +2,7 @@ package v1
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
@@ -32,9 +33,9 @@ func (h *BulkActionHandler) WithWorkerPool(wp *worker.WorkerPool) *BulkActionHan
|
||||
// BulkActionRequest is the DTO for generic bulk actions.
|
||||
// Reference: Chatwoot bulk_actions_controller#create — params[:type], params[:action_name], params[:ids]
|
||||
type BulkActionRequest struct {
|
||||
Type string `json:"type" binding:"required"` // "Conversation" or "Contact"
|
||||
ActionName string `json:"action_name,omitempty"` // legacy local action names
|
||||
IDs []uint `json:"ids" binding:"required,min=1"` // Chatwoot sends conversation display IDs
|
||||
Type string `json:"type"` // "Conversation" or "Contact"
|
||||
ActionName string `json:"action_name,omitempty"` // legacy local action names
|
||||
IDs []uint `json:"ids"` // Chatwoot sends conversation display IDs
|
||||
Fields service.ConversationBulkActionFields `json:"fields,omitempty"`
|
||||
Labels service.ConversationBulkActionLabels `json:"labels,omitempty"`
|
||||
AssigneeID *uint `json:"assignee_id,omitempty"` // legacy local assign field
|
||||
@@ -58,6 +59,7 @@ func (h *BulkActionHandler) Create(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
req.Type = normalizeBulkActionType(req.Type)
|
||||
switch req.Type {
|
||||
case "Conversation":
|
||||
h.handleConversationBulk(c, accountID, req)
|
||||
@@ -141,31 +143,80 @@ func (h *BulkActionHandler) handleConversationBulk(c *gin.Context, accountID uin
|
||||
// handleContactBulk processes bulk actions on contacts.
|
||||
// Reference: Chatwoot only supports "delete" and label operations for contacts in bulk_actions.
|
||||
func (h *BulkActionHandler) handleContactBulk(c *gin.Context, accountID uint, req BulkActionRequest) {
|
||||
successCount := 0
|
||||
failCount := 0
|
||||
|
||||
for _, contactID := range req.IDs {
|
||||
var svcErr error
|
||||
|
||||
switch req.ActionName {
|
||||
case "delete":
|
||||
svcErr = h.contactSvc.Delete(c.Request.Context(), accountID, contactID)
|
||||
default:
|
||||
// Chatwoot: other contact actions (labels) require separate async jobs
|
||||
// For now, return error for unsupported actions
|
||||
failCount++
|
||||
continue
|
||||
if h.worker != nil {
|
||||
params := service.ContactBulkActionParams{
|
||||
Type: req.Type,
|
||||
ActionName: req.ActionName,
|
||||
IDs: req.IDs,
|
||||
Labels: req.Labels,
|
||||
}
|
||||
|
||||
if svcErr != nil {
|
||||
failCount++
|
||||
} else {
|
||||
successCount++
|
||||
if _, err := service.EnqueueContactBulkAction(c.Request.Context(), h.worker, accountID, getUserID(c), params); err != nil {
|
||||
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to enqueue bulk action")
|
||||
return
|
||||
}
|
||||
c.Status(http.StatusOK)
|
||||
return
|
||||
}
|
||||
|
||||
response.OK(c, gin.H{
|
||||
"success_count": successCount,
|
||||
"fail_count": failCount,
|
||||
})
|
||||
if err := h.performContactBulkSync(c, accountID, req); err != nil {
|
||||
response.AbortWithStatusError(c, http.StatusInternalServerError, response.ErrInternal, "failed to process bulk action")
|
||||
return
|
||||
}
|
||||
c.Status(http.StatusOK)
|
||||
}
|
||||
|
||||
func (h *BulkActionHandler) performContactBulkSync(c *gin.Context, accountID uint, req BulkActionRequest) error {
|
||||
if h.contactSvc == nil {
|
||||
return nil
|
||||
}
|
||||
switch {
|
||||
case req.ActionName == "delete":
|
||||
for _, contactID := range req.IDs {
|
||||
if err := h.contactSvc.Delete(c.Request.Context(), accountID, contactID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
case len(req.Labels.Add) > 0:
|
||||
for _, contactID := range req.IDs {
|
||||
current, err := h.contactSvc.GetLabels(c.Request.Context(), accountID, contactID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := h.contactSvc.UpdateLabels(c.Request.Context(), accountID, contactID, append(current, req.Labels.Add...)); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
case len(req.Labels.Remove) > 0:
|
||||
remove := map[string]struct{}{}
|
||||
for _, label := range req.Labels.Remove {
|
||||
remove[strings.TrimSpace(label)] = struct{}{}
|
||||
}
|
||||
for _, contactID := range req.IDs {
|
||||
current, err := h.contactSvc.GetLabels(c.Request.Context(), accountID, contactID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
kept := current[:0]
|
||||
for _, label := range current {
|
||||
if _, ok := remove[label]; !ok {
|
||||
kept = append(kept, label)
|
||||
}
|
||||
}
|
||||
if _, err := h.contactSvc.UpdateLabels(c.Request.Context(), accountID, contactID, kept); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func normalizeBulkActionType(value string) string {
|
||||
switch strings.ToLower(strings.TrimSpace(value)) {
|
||||
case "conversation":
|
||||
return "Conversation"
|
||||
case "contact":
|
||||
return "Contact"
|
||||
default:
|
||||
return strings.TrimSpace(value)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -47,3 +47,50 @@ func TestBulkActionHandler_ConversationEnqueuesChatwootPayload(t *testing.T) {
|
||||
require.Contains(t, string(job.Payload), `"status":"resolved"`)
|
||||
require.Contains(t, string(job.Payload), `"add":["vip"]`)
|
||||
}
|
||||
|
||||
func TestBulkActionHandler_ContactEnqueuesChatwootPayload(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
db, err := gorm.Open(sqlite.Open("file:bulk-action-contact-handler?mode=memory&cache=private"), &gorm.Config{})
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, db.AutoMigrate(&model.BackgroundJob{}))
|
||||
sqlDB, _ := db.DB()
|
||||
t.Cleanup(func() { _ = sqlDB.Close() })
|
||||
|
||||
wp := worker.NewWorkerPool(db)
|
||||
handler := NewBulkActionHandler(nil, nil).WithWorkerPool(wp)
|
||||
router := gin.New()
|
||||
router.POST("/api/v1/accounts/:account_id/bulk_actions", handler.Create)
|
||||
|
||||
body := []byte(`{"type":"contact","ids":[11,12],"labels":{"add":["vip"],"remove":["old"]}}`)
|
||||
request := httptest.NewRequest(http.MethodPost, "/api/v1/accounts/7/bulk_actions", bytes.NewReader(body))
|
||||
request.Header.Set("Content-Type", "application/json")
|
||||
request.Header.Set("X-User-ID", "42")
|
||||
response := httptest.NewRecorder()
|
||||
router.ServeHTTP(response, request)
|
||||
|
||||
require.Equal(t, http.StatusOK, response.Code)
|
||||
require.Empty(t, response.Body.String())
|
||||
|
||||
var job model.BackgroundJob
|
||||
require.NoError(t, db.Where("job_type = ? AND queue = ?", service.TaskTypeContactBulkAction, "medium").First(&job).Error)
|
||||
require.Contains(t, string(job.Payload), `"account_id":7`)
|
||||
require.Contains(t, string(job.Payload), `"user_id":42`)
|
||||
require.Contains(t, string(job.Payload), `"type":"Contact"`)
|
||||
require.Contains(t, string(job.Payload), `"ids":[11,12]`)
|
||||
require.Contains(t, string(job.Payload), `"add":["vip"]`)
|
||||
}
|
||||
|
||||
func TestBulkActionHandler_InvalidTypeMatchesChatwootPayload(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
handler := NewBulkActionHandler(nil, nil)
|
||||
router := gin.New()
|
||||
router.POST("/api/v1/accounts/:account_id/bulk_actions", handler.Create)
|
||||
|
||||
request := httptest.NewRequest(http.MethodPost, "/api/v1/accounts/7/bulk_actions", bytes.NewReader([]byte(`{"type":"Ticket"}`)))
|
||||
request.Header.Set("Content-Type", "application/json")
|
||||
response := httptest.NewRecorder()
|
||||
router.ServeHTTP(response, request)
|
||||
|
||||
require.Equal(t, http.StatusUnprocessableEntity, response.Code)
|
||||
require.JSONEq(t, `{"success":false}`, response.Body.String())
|
||||
}
|
||||
|
||||
@@ -22,6 +22,7 @@ const (
|
||||
TaskTypeConversationResolutionForAccount = "conversation:resolution"
|
||||
TaskTypeConversationUpdateMessageStatus = "conversation:update_message_status"
|
||||
TaskTypeConversationBulkAction = "conversation:bulk_action"
|
||||
TaskTypeContactBulkAction = "contact:bulk_action"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -64,12 +65,25 @@ type ConversationBulkActionLabels struct {
|
||||
Remove []string `json:"remove,omitempty"`
|
||||
}
|
||||
|
||||
type ContactBulkActionParams struct {
|
||||
Type string `json:"type"`
|
||||
ActionName string `json:"action_name,omitempty"`
|
||||
IDs []uint `json:"ids"`
|
||||
Labels ConversationBulkActionLabels `json:"labels,omitempty"`
|
||||
}
|
||||
|
||||
type conversationBulkActionJob struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
UserID uint `json:"user_id,omitempty"`
|
||||
Params ConversationBulkActionParams `json:"params"`
|
||||
}
|
||||
|
||||
type contactBulkActionJob struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
UserID uint `json:"user_id,omitempty"`
|
||||
Params ContactBulkActionParams `json:"params"`
|
||||
}
|
||||
|
||||
var conversationMaintenanceRegistrations sync.Map
|
||||
|
||||
// RegisterConversationMaintenanceJobs wires Chatwoot scheduled maintenance jobs
|
||||
@@ -94,6 +108,7 @@ func registerConversationMaintenanceJobsWithNow(wp *worker.WorkerPool, db *gorm.
|
||||
wp.Register(TaskTypeConversationResolutionForAccount, runner.performResolutionForAccount)
|
||||
wp.Register(TaskTypeConversationUpdateMessageStatus, runner.performUpdateMessageStatus)
|
||||
wp.Register(TaskTypeConversationBulkAction, runner.performConversationBulkAction)
|
||||
wp.Register(TaskTypeContactBulkAction, runner.performContactBulkAction)
|
||||
}
|
||||
|
||||
func EnqueueScheduledItemsTrigger(ctx context.Context, wp *worker.WorkerPool, scheduledAt time.Time) (*model.BackgroundJob, error) {
|
||||
@@ -127,7 +142,7 @@ func EnqueueConversationBulkAction(ctx context.Context, wp *worker.WorkerPool, a
|
||||
if wp == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if accountID == 0 || len(params.IDs) == 0 {
|
||||
if accountID == 0 {
|
||||
return nil, fmt.Errorf("invalid conversation bulk action payload: account_id=%d ids=%v", accountID, params.IDs)
|
||||
}
|
||||
return wp.Enqueue(ctx, TaskTypeConversationBulkAction, conversationBulkActionJob{AccountID: accountID, UserID: userID, Params: params},
|
||||
@@ -136,6 +151,19 @@ func EnqueueConversationBulkAction(ctx context.Context, wp *worker.WorkerPool, a
|
||||
)
|
||||
}
|
||||
|
||||
func EnqueueContactBulkAction(ctx context.Context, wp *worker.WorkerPool, accountID, userID uint, params ContactBulkActionParams) (*model.BackgroundJob, error) {
|
||||
if wp == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if accountID == 0 {
|
||||
return nil, fmt.Errorf("invalid contact bulk action payload: account_id=%d ids=%v", accountID, params.IDs)
|
||||
}
|
||||
return wp.Enqueue(ctx, TaskTypeContactBulkAction, contactBulkActionJob{AccountID: accountID, UserID: userID, Params: params},
|
||||
worker.WithQueue("medium"),
|
||||
worker.WithMaxAttempts(3),
|
||||
)
|
||||
}
|
||||
|
||||
func scheduledItemsIdempotencyKey(scheduledAt time.Time) string {
|
||||
bucket := scheduledAt.UTC().Truncate(scheduledItemsInterval).Unix()
|
||||
return fmt.Sprintf("scheduled:trigger_items:%d", bucket)
|
||||
@@ -308,9 +336,12 @@ func (r *conversationMaintenanceRunner) performConversationBulkAction(ctx contex
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
return fmt.Errorf("unmarshal conversation bulk action job: %w", err)
|
||||
}
|
||||
if payload.AccountID == 0 || len(payload.Params.IDs) == 0 {
|
||||
if payload.AccountID == 0 {
|
||||
return fmt.Errorf("invalid conversation bulk action job payload: %#v", payload)
|
||||
}
|
||||
if len(payload.Params.IDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
if payload.Params.Type != "" && payload.Params.Type != "Conversation" {
|
||||
return nil
|
||||
}
|
||||
@@ -359,6 +390,85 @@ func (r *conversationMaintenanceRunner) performConversationBulkAction(ctx contex
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *conversationMaintenanceRunner) performContactBulkAction(ctx context.Context, job *model.BackgroundJob) error {
|
||||
var payload contactBulkActionJob
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
return fmt.Errorf("unmarshal contact bulk action job: %w", err)
|
||||
}
|
||||
if payload.AccountID == 0 {
|
||||
return fmt.Errorf("invalid contact bulk action job payload: %#v", payload)
|
||||
}
|
||||
if payload.Params.Type != "" && payload.Params.Type != "Contact" {
|
||||
return nil
|
||||
}
|
||||
if len(payload.Params.IDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
switch {
|
||||
case payload.Params.ActionName == "delete":
|
||||
return r.db.WithContext(ctx).
|
||||
Where("account_id = ? AND id IN ?", payload.AccountID, payload.Params.IDs).
|
||||
Delete(&model.Contact{}).Error
|
||||
case len(payload.Params.Labels.Add) > 0:
|
||||
return bulkAddContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Add)
|
||||
case len(payload.Params.Labels.Remove) > 0:
|
||||
return bulkRemoveContactLabels(ctx, r.db, payload.AccountID, payload.Params.IDs, payload.Params.Labels.Remove)
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func bulkAddContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) error {
|
||||
labels = normalizeContactServiceLabels(labels)
|
||||
if len(contactIDs) == 0 || len(labels) == 0 {
|
||||
return nil
|
||||
}
|
||||
var scopedContactIDs []uint
|
||||
if err := db.WithContext(ctx).Model(&model.Contact{}).
|
||||
Where("account_id = ? AND id IN ?", accountID, contactIDs).
|
||||
Pluck("id", &scopedContactIDs).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if len(scopedContactIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
for _, label := range labels {
|
||||
tag := model.Tag{AccountID: accountID, Name: label}
|
||||
if err := tx.Where("account_id = ? AND name = ?", accountID, label).FirstOrCreate(&tag).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
for _, contactID := range scopedContactIDs {
|
||||
contactLabel := model.ContactLabel{AccountID: accountID, ContactID: contactID, TagID: tag.ID}
|
||||
if err := tx.Where("account_id = ? AND contact_id = ? AND tag_id = ?", accountID, contactID, tag.ID).FirstOrCreate(&contactLabel).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func bulkRemoveContactLabels(ctx context.Context, db *gorm.DB, accountID uint, contactIDs []uint, labels []string) error {
|
||||
labels = normalizeContactServiceLabels(labels)
|
||||
if len(contactIDs) == 0 || len(labels) == 0 {
|
||||
return nil
|
||||
}
|
||||
var tagIDs []uint
|
||||
if err := db.WithContext(ctx).Model(&model.Tag{}).
|
||||
Where("account_id = ? AND name IN ?", accountID, labels).
|
||||
Pluck("id", &tagIDs).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if len(tagIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
return db.WithContext(ctx).
|
||||
Where("account_id = ? AND contact_id IN ? AND tag_id IN ?", accountID, contactIDs, tagIDs).
|
||||
Delete(&model.ContactLabel{}).Error
|
||||
}
|
||||
|
||||
func parseBulkActionTime(value string) (time.Time, bool) {
|
||||
if value == "" {
|
||||
return time.Time{}, false
|
||||
|
||||
@@ -318,6 +318,75 @@ func TestConversationMaintenanceJobsConversationBulkAction(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConversationMaintenanceJobsContactBulkAction(t *testing.T) {
|
||||
now := time.Date(2026, 6, 6, 0, 0, 0, 0, time.UTC)
|
||||
db := setupServiceTestDB(t)
|
||||
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return now }))
|
||||
registerConversationMaintenanceJobsWithNow(wp, db, func() time.Time { return now })
|
||||
|
||||
account := createTestAccount(t, db)
|
||||
otherAccount := createTestAccount(t, db)
|
||||
contactA := createTestContact(t, db, account.ID)
|
||||
contactB := createTestContact(t, db, account.ID)
|
||||
otherContact := createTestContact(t, db, otherAccount.ID)
|
||||
oldTag := model.Tag{AccountID: account.ID, Name: "old"}
|
||||
if err := db.Create(&oldTag).Error; err != nil {
|
||||
t.Fatalf("create old tag: %v", err)
|
||||
}
|
||||
if err := db.Create(&model.ContactLabel{AccountID: account.ID, ContactID: contactA.ID, TagID: oldTag.ID}).Error; err != nil {
|
||||
t.Fatalf("label contact A: %v", err)
|
||||
}
|
||||
if err := db.Create(&model.ContactLabel{AccountID: account.ID, ContactID: contactB.ID, TagID: oldTag.ID}).Error; err != nil {
|
||||
t.Fatalf("label contact B: %v", err)
|
||||
}
|
||||
|
||||
_, err := EnqueueContactBulkAction(context.Background(), wp, account.ID, 42, ContactBulkActionParams{
|
||||
Type: "Contact",
|
||||
IDs: []uint{contactA.ID, contactB.ID, otherContact.ID},
|
||||
Labels: ConversationBulkActionLabels{Add: []string{"vip", "vip"}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue contact label add: %v", err)
|
||||
}
|
||||
assertJobCount(t, db, TaskTypeContactBulkAction, 1)
|
||||
processRequiredJob(t, wp, "contact bulk add labels")
|
||||
assertContactHasLabel(t, db, account.ID, contactA.ID, "vip", true)
|
||||
assertContactHasLabel(t, db, account.ID, contactB.ID, "vip", true)
|
||||
assertContactHasLabel(t, db, otherAccount.ID, otherContact.ID, "vip", false)
|
||||
|
||||
_, err = EnqueueContactBulkAction(context.Background(), wp, account.ID, 42, ContactBulkActionParams{
|
||||
Type: "Contact",
|
||||
IDs: []uint{contactA.ID, contactB.ID},
|
||||
Labels: ConversationBulkActionLabels{Remove: []string{"old"}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue contact label remove: %v", err)
|
||||
}
|
||||
processRequiredJob(t, wp, "contact bulk remove labels")
|
||||
assertContactHasLabel(t, db, account.ID, contactA.ID, "old", false)
|
||||
assertContactHasLabel(t, db, account.ID, contactB.ID, "old", false)
|
||||
assertContactHasLabel(t, db, account.ID, contactA.ID, "vip", true)
|
||||
|
||||
_, err = EnqueueContactBulkAction(context.Background(), wp, account.ID, 42, ContactBulkActionParams{
|
||||
Type: "Contact",
|
||||
ActionName: "delete",
|
||||
IDs: []uint{contactA.ID, otherContact.ID},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("enqueue contact delete: %v", err)
|
||||
}
|
||||
processRequiredJob(t, wp, "contact bulk delete")
|
||||
|
||||
var deleted model.Contact
|
||||
if err := db.First(&deleted, contactA.ID).Error; err == nil {
|
||||
t.Fatalf("expected account contact to be soft deleted")
|
||||
}
|
||||
var untouched model.Contact
|
||||
if err := db.First(&untouched, otherContact.ID).Error; err != nil {
|
||||
t.Fatalf("other account contact should remain visible: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func createTestOneoffCampaign(t *testing.T, db *gorm.DB, accountID, inboxID, contactID uint, scheduledAt time.Time) *campaign.Campaign {
|
||||
t.Helper()
|
||||
c := &campaign.Campaign{
|
||||
@@ -358,6 +427,20 @@ func createConversationMaintenanceMessage(t *testing.T, db *gorm.DB, accountID,
|
||||
return message
|
||||
}
|
||||
|
||||
func assertContactHasLabel(t *testing.T, db *gorm.DB, accountID, contactID uint, label string, want bool) {
|
||||
t.Helper()
|
||||
var count int64
|
||||
if err := db.Table("contact_labels").
|
||||
Joins("JOIN tags ON tags.id = contact_labels.tag_id").
|
||||
Where("contact_labels.account_id = ? AND contact_labels.contact_id = ? AND tags.name = ?", accountID, contactID, label).
|
||||
Count(&count).Error; err != nil {
|
||||
t.Fatalf("count contact label %s: %v", label, err)
|
||||
}
|
||||
if got := count > 0; got != want {
|
||||
t.Fatalf("contact %d label %q presence = %v, want %v", contactID, label, got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func assertMessageStatus(t *testing.T, db *gorm.DB, messageID uint, want string) {
|
||||
t.Helper()
|
||||
var message model.Message
|
||||
|
||||
Reference in New Issue
Block a user