feat(conversations): queue message status updates
This commit is contained in:
@@ -694,6 +694,7 @@ func Bootstrap(env string) (*App, error) {
|
||||
widgetFileUploadRepo := repository.NewWidgetFileUploadRepo(db)
|
||||
widgetOfflineMessageRepo := repository.NewWidgetOfflineMessageRepo(db)
|
||||
widgetService := service.NewWidgetService(inboxRepo, contactRepo, contactInboxRepo, conversationRepo, messageRepo, widgetTypingAdapter, widgetThemeConfigRepo, preChatFormRepo, widgetFileUploadRepo, widgetOfflineMessageRepo, inboxMemberRepo, tagRepo, campaignRepo)
|
||||
widgetService.SetWorkerPool(workerPool)
|
||||
widgetHandler := widget.NewHandler(widgetService)
|
||||
|
||||
// Upload: DirectUpload repo + service + handler (account-level + widget direct uploads)
|
||||
|
||||
@@ -19,6 +19,7 @@ const (
|
||||
TaskTypeConversationReopenSnoozed = "conversation:reopen_snoozed"
|
||||
TaskTypeConversationResolutionScheduler = "account:conversations_resolution_scheduler"
|
||||
TaskTypeConversationResolutionForAccount = "conversation:resolution"
|
||||
TaskTypeConversationUpdateMessageStatus = "conversation:update_message_status"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -35,6 +36,12 @@ type conversationResolutionJob struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
}
|
||||
|
||||
type conversationUpdateMessageStatusJob struct {
|
||||
ConversationID uint `json:"conversation_id"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
Status string `json:"status"`
|
||||
}
|
||||
|
||||
var conversationMaintenanceRegistrations sync.Map
|
||||
|
||||
// RegisterConversationMaintenanceJobs wires Chatwoot scheduled maintenance jobs
|
||||
@@ -57,6 +64,7 @@ func registerConversationMaintenanceJobsWithNow(wp *worker.WorkerPool, db *gorm.
|
||||
wp.Register(TaskTypeConversationReopenSnoozed, runner.performReopenSnoozed)
|
||||
wp.Register(TaskTypeConversationResolutionScheduler, runner.performResolutionScheduler)
|
||||
wp.Register(TaskTypeConversationResolutionForAccount, runner.performResolutionForAccount)
|
||||
wp.Register(TaskTypeConversationUpdateMessageStatus, runner.performUpdateMessageStatus)
|
||||
}
|
||||
|
||||
func EnqueueScheduledItemsTrigger(ctx context.Context, wp *worker.WorkerPool, scheduledAt time.Time) (*model.BackgroundJob, error) {
|
||||
@@ -71,6 +79,21 @@ func EnqueueScheduledItemsTrigger(ctx context.Context, wp *worker.WorkerPool, sc
|
||||
)
|
||||
}
|
||||
|
||||
func EnqueueConversationMessageStatusUpdate(ctx context.Context, wp *worker.WorkerPool, conversationID uint, timestamp time.Time, status string) (*model.BackgroundJob, error) {
|
||||
if wp == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if status == "" {
|
||||
status = string(model.MessageStatusRead)
|
||||
}
|
||||
payload := conversationUpdateMessageStatusJob{ConversationID: conversationID, Timestamp: timestamp.UTC().Unix(), Status: status}
|
||||
return wp.Enqueue(ctx, TaskTypeConversationUpdateMessageStatus, payload,
|
||||
worker.WithQueue("deferred"),
|
||||
worker.WithMaxAttempts(3),
|
||||
worker.WithIdempotencyKey(fmt.Sprintf("conversation:update_message_status:%d:%s:%d", conversationID, status, payload.Timestamp)),
|
||||
)
|
||||
}
|
||||
|
||||
func scheduledItemsIdempotencyKey(scheduledAt time.Time) string {
|
||||
bucket := scheduledAt.UTC().Truncate(scheduledItemsInterval).Unix()
|
||||
return fmt.Sprintf("scheduled:trigger_items:%d", bucket)
|
||||
@@ -200,3 +223,40 @@ func (r *conversationMaintenanceRunner) performResolutionForAccount(ctx context.
|
||||
Limit(conversationResolutionLimit).
|
||||
Updates(updates).Error
|
||||
}
|
||||
|
||||
func (r *conversationMaintenanceRunner) performUpdateMessageStatus(ctx context.Context, job *model.BackgroundJob) error {
|
||||
var payload conversationUpdateMessageStatusJob
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
return fmt.Errorf("unmarshal conversation message status job: %w", err)
|
||||
}
|
||||
if payload.ConversationID == 0 || payload.Timestamp == 0 {
|
||||
return fmt.Errorf("invalid conversation message status job payload: %#v", payload)
|
||||
}
|
||||
if !validConversationMessageStatus(payload.Status) {
|
||||
return nil
|
||||
}
|
||||
|
||||
var conversation model.Conversation
|
||||
if err := r.db.WithContext(ctx).First(&conversation, payload.ConversationID).Error; err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("load conversation %d for message status update: %w", payload.ConversationID, err)
|
||||
}
|
||||
|
||||
return r.db.WithContext(ctx).Model(&model.Message{}).
|
||||
Where("conversation_id = ?", conversation.ID).
|
||||
Where("status IN ?", []string{string(model.MessageStatusSent), string(model.MessageStatusDelivered)}).
|
||||
Where("message_type <> ?", "incoming").
|
||||
Where("created_at <= ?", time.Unix(payload.Timestamp, 0).UTC()).
|
||||
Update("status", payload.Status).Error
|
||||
}
|
||||
|
||||
func validConversationMessageStatus(status string) bool {
|
||||
switch status {
|
||||
case string(model.MessageStatusRead), string(model.MessageStatusDelivered):
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -179,6 +179,64 @@ func TestConversationMaintenanceJobsRetryMissingResolutionAccount(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestConversationMaintenanceJobsUpdateMessageStatus(t *testing.T) {
|
||||
now := time.Date(2026, 6, 5, 22, 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)
|
||||
inbox := createTestInbox(t, db, account.ID, "web_widget")
|
||||
contact := createTestContact(t, db, account.ID)
|
||||
conversation := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
|
||||
cutoff := now.Add(-time.Minute)
|
||||
|
||||
beforeSent := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "outgoing", string(model.MessageStatusSent), cutoff.Add(-time.Minute))
|
||||
beforeDelivered := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "outgoing", string(model.MessageStatusDelivered), cutoff.Add(-30*time.Second))
|
||||
incoming := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "incoming", string(model.MessageStatusSent), cutoff.Add(-time.Minute))
|
||||
alreadyRead := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "outgoing", string(model.MessageStatusRead), cutoff.Add(-time.Minute))
|
||||
afterCutoff := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "outgoing", string(model.MessageStatusSent), cutoff.Add(time.Minute))
|
||||
|
||||
if _, err := EnqueueConversationMessageStatusUpdate(context.Background(), wp, conversation.ID, cutoff, string(model.MessageStatusRead)); err != nil {
|
||||
t.Fatalf("enqueue message status update: %v", err)
|
||||
}
|
||||
assertJobCount(t, db, TaskTypeConversationUpdateMessageStatus, 1)
|
||||
var queued model.BackgroundJob
|
||||
if err := db.Where("job_type = ?", TaskTypeConversationUpdateMessageStatus).First(&queued).Error; err != nil {
|
||||
t.Fatalf("load queued message status job: %v", err)
|
||||
}
|
||||
if queued.Queue != "deferred" {
|
||||
t.Fatalf("expected deferred queue, got %s", queued.Queue)
|
||||
}
|
||||
|
||||
processRequiredJob(t, wp, "message status")
|
||||
|
||||
assertMessageStatus(t, db, beforeSent.ID, string(model.MessageStatusRead))
|
||||
assertMessageStatus(t, db, beforeDelivered.ID, string(model.MessageStatusRead))
|
||||
assertMessageStatus(t, db, incoming.ID, string(model.MessageStatusSent))
|
||||
assertMessageStatus(t, db, alreadyRead.ID, string(model.MessageStatusRead))
|
||||
assertMessageStatus(t, db, afterCutoff.ID, string(model.MessageStatusSent))
|
||||
}
|
||||
|
||||
func TestConversationMaintenanceJobsIgnoreInvalidMessageStatus(t *testing.T) {
|
||||
now := time.Date(2026, 6, 5, 22, 30, 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)
|
||||
inbox := createTestInbox(t, db, account.ID, "web_widget")
|
||||
contact := createTestContact(t, db, account.ID)
|
||||
conversation := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
|
||||
message := createConversationMaintenanceMessage(t, db, account.ID, inbox.ID, conversation.ID, "outgoing", string(model.MessageStatusSent), now.Add(-time.Minute))
|
||||
|
||||
if _, err := wp.Enqueue(context.Background(), TaskTypeConversationUpdateMessageStatus, conversationUpdateMessageStatusJob{ConversationID: conversation.ID, Timestamp: now.Unix(), Status: "failed"}, worker.WithQueue("deferred")); err != nil {
|
||||
t.Fatalf("enqueue invalid status job: %v", err)
|
||||
}
|
||||
processRequiredJob(t, wp, "invalid message status")
|
||||
assertMessageStatus(t, db, message.ID, string(model.MessageStatusSent))
|
||||
}
|
||||
|
||||
func createTestOneoffCampaign(t *testing.T, db *gorm.DB, accountID, inboxID, contactID uint, scheduledAt time.Time) *campaign.Campaign {
|
||||
t.Helper()
|
||||
c := &campaign.Campaign{
|
||||
@@ -201,6 +259,35 @@ func createTestOneoffCampaign(t *testing.T, db *gorm.DB, accountID, inboxID, con
|
||||
return c
|
||||
}
|
||||
|
||||
func createConversationMaintenanceMessage(t *testing.T, db *gorm.DB, accountID, inboxID, conversationID uint, messageType, status string, createdAt time.Time) *model.Message {
|
||||
t.Helper()
|
||||
message := &model.Message{
|
||||
AccountID: accountID,
|
||||
InboxID: inboxID,
|
||||
ConversationID: conversationID,
|
||||
Content: fmt.Sprintf("%s %d", messageType, time.Now().UnixNano()),
|
||||
MessageType: messageType,
|
||||
Status: status,
|
||||
}
|
||||
message.CreatedAt = createdAt
|
||||
message.UpdatedAt = createdAt
|
||||
if err := db.Create(message).Error; err != nil {
|
||||
t.Fatalf("create message: %v", err)
|
||||
}
|
||||
return message
|
||||
}
|
||||
|
||||
func assertMessageStatus(t *testing.T, db *gorm.DB, messageID uint, want string) {
|
||||
t.Helper()
|
||||
var message model.Message
|
||||
if err := db.First(&message, messageID).Error; err != nil {
|
||||
t.Fatalf("load message %d: %v", messageID, err)
|
||||
}
|
||||
if message.Status != want {
|
||||
t.Fatalf("expected message %d status %s, got %s", messageID, want, message.Status)
|
||||
}
|
||||
}
|
||||
|
||||
func assertJobCount(t *testing.T, db *gorm.DB, jobType string, want int64) {
|
||||
t.Helper()
|
||||
var count int64
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
channelmodel "github.com/gochat/gochat/internal/model/channel"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/search"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
ws "github.com/gochat/gochat/internal/ws"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
"gorm.io/datatypes"
|
||||
@@ -51,6 +52,7 @@ type WidgetService struct {
|
||||
inboxMemberRepo *repository.InboxMemberRepo
|
||||
tagRepo *repository.TagRepo
|
||||
campaignRepo *repository.CampaignRepo
|
||||
worker *worker.WorkerPool
|
||||
}
|
||||
|
||||
// NewWidgetService creates a new Widget service.
|
||||
@@ -86,6 +88,10 @@ func NewWidgetService(
|
||||
}
|
||||
}
|
||||
|
||||
func (s *WidgetService) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
s.worker = wp
|
||||
}
|
||||
|
||||
// --- DTOs ---
|
||||
|
||||
// WidgetInitRequest is the DTO for the /widget/init endpoint.
|
||||
@@ -744,6 +750,9 @@ func (s *WidgetService) PublicUpdateLastSeen(ctx context.Context, inboxIdentifie
|
||||
if err := s.conversationRepo.Update(ctx, conversation); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := EnqueueConversationMessageStatusUpdate(ctx, s.worker, conversation.ID, time.Unix(now, 0), string(model.MessageStatusRead)); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return conversation, nil
|
||||
}
|
||||
|
||||
@@ -1109,6 +1118,9 @@ func (s *WidgetService) UpdateLastSeen(ctx context.Context, widgetToken string)
|
||||
if err := s.conversationRepo.Update(ctx, conversation); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := EnqueueConversationMessageStatusUpdate(ctx, s.worker, conversation.ID, time.Unix(now, 0), string(model.MessageStatusRead)); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return conversation, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
)
|
||||
|
||||
// ========== Test Helpers ==========
|
||||
@@ -44,6 +45,7 @@ func setupWidgetServiceTest(t *testing.T) (*gorm.DB, *WidgetService) {
|
||||
&model.PreChatForm{},
|
||||
&model.WidgetFileUpload{},
|
||||
&model.WidgetOfflineMessage{},
|
||||
&model.BackgroundJob{},
|
||||
), "failed to auto-migrate models")
|
||||
|
||||
t.Cleanup(func() {
|
||||
@@ -435,6 +437,41 @@ func TestWidgetService_GetConversations_InvalidToken(t *testing.T) {
|
||||
assert.Error(t, err)
|
||||
}
|
||||
|
||||
func TestWidgetService_UpdateLastSeenQueuesMessageStatusJob(t *testing.T) {
|
||||
db, svc := setupWidgetServiceTest(t)
|
||||
ctx := context.Background()
|
||||
|
||||
seedWidgetInbox(t, db)
|
||||
initResp, err := svc.Init(ctx, WidgetInitRequest{WebsiteToken: "test_ws_token_123"})
|
||||
require.NoError(t, err)
|
||||
|
||||
incoming, err := svc.SendMessage(ctx, WidgetSendMessageRequest{WidgetToken: initResp.WidgetToken, Content: "hello"})
|
||||
require.NoError(t, err)
|
||||
outgoing := &model.Message{
|
||||
AccountID: incoming.Message.AccountID,
|
||||
InboxID: incoming.Message.InboxID,
|
||||
ConversationID: incoming.ConversationID,
|
||||
Content: "reply before last seen",
|
||||
MessageType: "outgoing",
|
||||
Status: string(model.MessageStatusSent),
|
||||
}
|
||||
outgoing.CreatedAt = time.Now().Add(-time.Minute)
|
||||
outgoing.UpdatedAt = outgoing.CreatedAt
|
||||
require.NoError(t, db.Create(outgoing).Error)
|
||||
|
||||
wp := worker.NewWorkerPool(db)
|
||||
RegisterConversationMaintenanceJobs(wp, db)
|
||||
svc.SetWorkerPool(wp)
|
||||
|
||||
conversation, err := svc.UpdateLastSeen(ctx, initResp.WidgetToken)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, conversation.ContactLastSeenAt)
|
||||
|
||||
var job model.BackgroundJob
|
||||
require.NoError(t, db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeConversationUpdateMessageStatus, "deferred", model.BackgroundJobStatusQueued).First(&job).Error)
|
||||
assert.Contains(t, string(job.Payload), fmt.Sprintf(`"conversation_id":%d`, outgoing.ConversationID))
|
||||
}
|
||||
|
||||
// ========== GetMessages Tests ==========
|
||||
|
||||
func TestWidgetService_GetMessages(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user