feat(conversations): align destroy job

This commit is contained in:
2026-06-07 12:26:48 +08:00
parent bec65a86e4
commit 27f927327d
8 changed files with 152 additions and 7 deletions
+1
View File
@@ -503,6 +503,7 @@ func Bootstrap(env string) (*App, error) {
accountUserRepo := repository.NewAccountUserRepo(db)
agentRepo := repository.NewAgentRepo(db)
conversationService := service.NewConversationService(conversationRepo, messageRepo, channelDispatcher, inboxMemberService, accountUserRepo, teamRepo, teamMemberRepo)
conversationService.SetWorkerPool(workerPool)
appliedSlaService := service.NewAppliedSlaService(appliedSlaRepo, slaEventRepo, slaPolicyRepo, conversationRepo)
conversationService.SetAppliedSlaService(appliedSlaService)
service.RegisterSlaProcessingJobs(workerPool, db, appliedSlaService)
@@ -212,7 +212,7 @@ func (h *ConversationHandler) Delete(c *gin.Context) {
Action: "destroy",
AuditedChanges: gin.H{"id": conversation.ID, "display_id": conversation.DisplayID},
})
response.NoContent(c)
c.Status(http.StatusOK)
}
// @Summary Assign an agent to a conversation
@@ -589,7 +589,8 @@ func (s *ConversationCrudTestSuite) TestDelete_Success() {
req, _ := http.NewRequest("DELETE", s.convURL(s.testConv.ID), nil)
s.router.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusNoContent, w.Code)
assert.Equal(s.T(), http.StatusOK, w.Code)
assert.Empty(s.T(), w.Body.String())
}
func (s *ConversationCrudTestSuite) TestDelete_InvalidAccountID() {
@@ -144,6 +144,7 @@ func (s *ConversationHandlerTestSuite) SetupSuite() {
conversations.GET("/:conversation_id/reporting_events", handler.ReportingEvents)
conversations.POST("/:conversation_id/toggle_typing", handler.ToggleTyping)
conversations.POST("/:conversation_id/update_last_seen", handler.UpdateLastSeen)
conversations.DELETE("/:conversation_id", handler.Delete)
}
}
}
@@ -448,6 +449,26 @@ func (s *ConversationHandlerTestSuite) TestTranscript_ConversationNotFound() {
assert.Equal(s.T(), http.StatusNotFound, w.Code)
}
func (s *ConversationHandlerTestSuite) TestDelete_SuccessReturnsChatwootHeadOK() {
w := httptest.NewRecorder()
req, _ := http.NewRequest(http.MethodDelete, s.accountURL()+"/conversations/"+strconv.FormatUint(uint64(s.testConv.ID), 10), nil)
s.router.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusOK, w.Code)
assert.Empty(s.T(), w.Body.String())
var deleted model.Conversation
assert.Error(s.T(), s.db.First(&deleted, s.testConv.ID).Error)
}
func (s *ConversationHandlerTestSuite) TestDelete_ConversationNotFound() {
w := httptest.NewRecorder()
req, _ := http.NewRequest(http.MethodDelete, s.accountURL()+"/conversations/9999", nil)
s.router.ServeHTTP(w, req)
assert.Equal(s.T(), http.StatusNotFound, w.Code)
}
// ========== UpdateCustomAttributes Handler Tests ==========
func (s *ConversationHandlerTestSuite) TestUpdateCustomAttributes_Success() {
@@ -0,0 +1,40 @@
package service
import (
"context"
"encoding/json"
"sync"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
const TaskTypeConversationDeleteObject = "conversation:delete_object"
type conversationDeleteObjectJob struct {
AccountID uint `json:"account_id"`
ConversationID uint `json:"conversation_id"`
}
var conversationDeleteRegistrations sync.Map
// RegisterConversationDeleteJobs wires Chatwoot's DeleteObjectJob path for
// conversation destroy actions. Without a worker, ConversationService.Delete
// keeps the synchronous fallback used by focused tests.
func RegisterConversationDeleteJobs(wp *worker.WorkerPool, svc *ConversationService) {
if wp == nil || svc == nil {
return
}
if _, loaded := conversationDeleteRegistrations.LoadOrStore(wp, struct{}{}); loaded {
return
}
wp.Register(TaskTypeConversationDeleteObject, svc.performConversationDeleteObject)
}
func (s *ConversationService) performConversationDeleteObject(ctx context.Context, job *model.BackgroundJob) error {
var payload conversationDeleteObjectJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return err
}
return s.deleteNow(ctx, payload.AccountID, payload.ConversationID)
}
+23 -1
View File
@@ -13,6 +13,7 @@ import (
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/gochat/gochat/internal/search"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
pkgvalidator "github.com/gochat/gochat/pkg/validator"
@@ -33,6 +34,7 @@ type ConversationService struct {
searchIndexer SearchIndexer
appliedSlaSvc *AppliedSlaService
transcriptMailer automation.AutomationTranscriptDeliverer
worker *worker.WorkerPool
}
// NewConversationService creates a new Conversation service.
@@ -52,6 +54,11 @@ func (s *ConversationService) SetTranscriptDeliverer(deliverer automation.Automa
s.transcriptMailer = deliverer
}
func (s *ConversationService) SetWorkerPool(wp *worker.WorkerPool) {
s.worker = wp
RegisterConversationDeleteJobs(wp, s)
}
func (s *ConversationService) DB() *gorm.DB {
if s == nil || s.repo == nil {
return nil
@@ -629,14 +636,29 @@ func (s *ConversationService) Delete(ctx context.Context, accountID, id uint) er
if err != nil {
return err
}
if s.worker != nil {
_, err := s.worker.Enqueue(ctx, TaskTypeConversationDeleteObject, conversationDeleteObjectJob{AccountID: accountID, ConversationID: conversation.ID}, worker.WithQueue("low"), worker.WithMaxAttempts(3), worker.WithIdempotencyKey(fmt.Sprintf("conversation-delete:%d:%d", accountID, conversation.ID)))
return err
}
return s.deleteLoaded(ctx, conversation)
}
func (s *ConversationService) deleteNow(ctx context.Context, accountID, id uint) error {
conversation, err := s.repo.FindByAccountAndID(ctx, accountID, id)
if err != nil {
return err
}
return s.deleteLoaded(ctx, conversation)
}
func (s *ConversationService) deleteLoaded(ctx context.Context, conversation *model.Conversation) error {
if err := s.repo.Delete(ctx, conversation.ID); err != nil {
return err
}
// Dispatch EventConversationDeleted
s.dispatchConversationEvent(ctx, channel.EventConversationDeleted, conversation)
s.deleteConversationIndex(ctx, accountID, conversation.ID)
s.deleteConversationIndex(ctx, conversation.AccountID, conversation.ID)
return nil
}
@@ -16,6 +16,7 @@ import (
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/gochat/gochat/internal/worker"
)
// ========== Test Setup ==========
@@ -48,6 +49,7 @@ func setupConversationServiceTestDB(t *testing.T) *gorm.DB {
&model.SlaPolicy{},
&model.AppliedSLA{},
&model.SlaEvent{},
&model.BackgroundJob{},
), "failed to auto-migrate")
t.Cleanup(func() {
@@ -779,6 +781,60 @@ func TestConversationService_SendTranscript_RateLimited(t *testing.T) {
assert.Empty(t, deliverer.requests)
}
func TestConversationService_Delete(t *testing.T) {
svc, db := setupConversationService(t)
account := createConversationServiceTestAccount(t, db)
inbox := createConversationServiceTestInbox(t, db, account.ID)
contact := createConversationServiceTestContact(t, db, account.ID)
conv := createConversationServiceTestConversation(t, db, account.ID, inbox.ID, contact.ID, "open")
require.NoError(t, svc.Delete(context.Background(), account.ID, conv.ID))
var found model.Conversation
assert.Error(t, db.First(&found, conv.ID).Error)
require.NoError(t, db.Unscoped().First(&found, conv.ID).Error)
assert.True(t, found.DeletedAt.Valid)
}
func TestConversationService_Delete_QueuesDeleteObjectJob(t *testing.T) {
svc, db := setupConversationService(t)
wp := worker.NewWorkerPool(db)
svc.SetWorkerPool(wp)
account := createConversationServiceTestAccount(t, db)
inbox := createConversationServiceTestInbox(t, db, account.ID)
contact := createConversationServiceTestContact(t, db, account.ID)
conv := createConversationServiceTestConversation(t, db, account.ID, inbox.ID, contact.ID, "open")
require.NoError(t, svc.Delete(context.Background(), account.ID, conv.ID))
var queued model.BackgroundJob
require.NoError(t, db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeConversationDeleteObject, "low", model.BackgroundJobStatusQueued).First(&queued).Error)
var stillVisible model.Conversation
require.NoError(t, db.First(&stillVisible, conv.ID).Error)
processed, err := wp.ProcessOne(context.Background())
require.NoError(t, err)
assert.True(t, processed)
var deleted model.Conversation
assert.Error(t, db.First(&deleted, conv.ID).Error)
}
func TestConversationService_DeleteWithWorker_NotFoundDoesNotQueue(t *testing.T) {
svc, db := setupConversationService(t)
wp := worker.NewWorkerPool(db)
svc.SetWorkerPool(wp)
err := svc.Delete(context.Background(), 1, 9999)
require.ErrorIs(t, err, gorm.ErrRecordNotFound)
var count int64
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeConversationDeleteObject).Count(&count).Error)
assert.Zero(t, count)
}
// ========== UpdateCustomAttributes Tests ==========
func TestConversationService_UpdateCustomAttributes(t *testing.T) {