feat(messages): align retry status parity
This commit is contained in:
@@ -13,6 +13,7 @@ import (
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/suite"
|
||||
"gorm.io/datatypes"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
@@ -560,6 +561,11 @@ func (s *MessageHandlerTestSuite) TestDelete_NotFound() {
|
||||
// --- Retry Tests ---
|
||||
|
||||
func (s *MessageHandlerTestSuite) TestRetry_Success() {
|
||||
s.Require().NoError(s.db.Model(s.testMessage).Updates(map[string]interface{}{
|
||||
"status": "failed",
|
||||
"content_attributes": datatypes.JSON([]byte(`{"external_error":"provider failed"}`)),
|
||||
}).Error)
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
url := msgRetryURL(s.testAccount.ID, s.testConv.ID, s.testMessage.ID)
|
||||
req, _ := http.NewRequest("POST", url, nil)
|
||||
@@ -570,6 +576,8 @@ func (s *MessageHandlerTestSuite) TestRetry_Success() {
|
||||
json.Unmarshal(w.Body.Bytes(), &resp)
|
||||
assert.Equal(s.T(), float64(s.testMessage.ID), resp["id"])
|
||||
assert.Equal(s.T(), float64(1), resp["message_type"])
|
||||
assert.Equal(s.T(), "sent", resp["status"])
|
||||
assert.Equal(s.T(), map[string]interface{}{}, resp["content_attributes"])
|
||||
}
|
||||
|
||||
func (s *MessageHandlerTestSuite) TestRetry_InvalidAccountID() {
|
||||
|
||||
@@ -438,27 +438,33 @@ func attachmentThumbURL(contentType string, messageID uint, fileName string) str
|
||||
}
|
||||
|
||||
// Retry retries a failed message by resetting its delivery status.
|
||||
// Reference: Chatwoot MessagesController#retry sets status to sent, clears
|
||||
// content_attributes, and queues SendReplyJob.
|
||||
func (s *MessageService) Retry(ctx context.Context, accountID, id uint) (*model.Message, error) {
|
||||
message, err := s.repo.FindByAccountAndID(ctx, accountID, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Reset status to indicate retry
|
||||
message.Status = "retrying"
|
||||
message.Status = "sent"
|
||||
message.ContentAttributes = datatypes.JSON([]byte(`{}`))
|
||||
|
||||
if err := s.repo.Update(ctx, message); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if s.worker != nil {
|
||||
if _, err := EnqueueSendReply(ctx, s.worker, message.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// Dispatch EventMessageStatusUpdated with "retrying" status
|
||||
event := channel.NewChannelEvent(channel.EventMessageStatusUpdated, channel.ChannelAPI, message.AccountID, message.InboxID)
|
||||
event.ConversationID = message.ConversationID
|
||||
if message.SenderID != nil {
|
||||
event.UserID = *message.SenderID
|
||||
}
|
||||
event.Data["message_id"] = message.ID
|
||||
event.Data["status"] = "retrying"
|
||||
event.Data["status"] = "sent"
|
||||
applogger.L().Infof("dispatching retry event for message %d", message.ID)
|
||||
if err := s.dispatcher.Dispatch(ctx, event); err != nil {
|
||||
applogger.L().Errorf("failed to dispatch retry event for message %d: %v", message.ID, err)
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/datatypes"
|
||||
"gorm.io/gorm"
|
||||
|
||||
"github.com/gochat/gochat/internal/channel"
|
||||
@@ -14,6 +15,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"
|
||||
)
|
||||
|
||||
// mockMessageLLMProvider implements llm.Provider for message service testing.
|
||||
@@ -579,7 +581,7 @@ func TestMessageService_Retry(t *testing.T) {
|
||||
msg := &model.Message{
|
||||
ConversationID: conv.ID, AccountID: account.ID, InboxID: inbox.ID,
|
||||
Content: "failed message", MessageType: "outgoing", ContentType: "text", SenderType: "user",
|
||||
Status: "failed",
|
||||
Status: "failed", ContentAttributes: datatypes.JSON([]byte(`{"external_error":"provider failed"}`)),
|
||||
}
|
||||
require.NoError(t, db.Create(msg).Error)
|
||||
|
||||
@@ -589,10 +591,43 @@ func TestMessageService_Retry(t *testing.T) {
|
||||
|
||||
retried, err := svc.Retry(ctx, account.ID, msg.ID)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, "retrying", retried.Status)
|
||||
assert.Equal(t, "sent", retried.Status)
|
||||
assert.JSONEq(t, `{}`, string(retried.ContentAttributes))
|
||||
assert.True(t, listener.received)
|
||||
assert.Equal(t, "retrying", listener.lastData["status"])
|
||||
assert.Equal(t, "sent", listener.lastData["status"])
|
||||
assert.Equal(t, msg.ID, listener.lastData["message_id"])
|
||||
|
||||
var stored model.Message
|
||||
require.NoError(t, db.First(&stored, msg.ID).Error)
|
||||
assert.JSONEq(t, `{}`, string(stored.ContentAttributes))
|
||||
})
|
||||
|
||||
t.Run("queues_send_reply_when_worker_configured", func(t *testing.T) {
|
||||
db, _, _, svc := setupMessageServiceWithDefaultLLM(t)
|
||||
wp := worker.NewWorkerPool(db)
|
||||
svc.SetWorkerPool(wp)
|
||||
|
||||
account := createTestAccount(t, db)
|
||||
inbox := createTestInbox(t, db, account.ID, "web_widget")
|
||||
contact := createTestContact(t, db, account.ID)
|
||||
conv := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
|
||||
|
||||
msg := &model.Message{
|
||||
ConversationID: conv.ID, AccountID: account.ID, InboxID: inbox.ID,
|
||||
Content: "failed message", MessageType: "outgoing", ContentType: "text", SenderType: "user",
|
||||
Status: "failed", ContentAttributes: datatypes.JSON([]byte(`{"external_error":"provider failed"}`)),
|
||||
}
|
||||
require.NoError(t, db.Create(msg).Error)
|
||||
|
||||
retried, err := svc.Retry(ctx, account.ID, msg.ID)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "sent", retried.Status)
|
||||
|
||||
var job model.BackgroundJob
|
||||
require.NoError(t, db.Where("job_type = ?", TaskTypeMessageSendReply).First(&job).Error)
|
||||
assert.Equal(t, model.BackgroundJobStatusQueued, job.Status)
|
||||
assert.Equal(t, "message:send_reply", job.JobType)
|
||||
assert.JSONEq(t, fmt.Sprintf(`{"message_id":%d}`, msg.ID), string(job.Payload))
|
||||
})
|
||||
|
||||
t.Run("message_not_found", func(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user