Files
gochat/internal/service/message_delivery_worker_test.go
T

213 lines
7.8 KiB
Go

package service
import (
"context"
"encoding/json"
"errors"
"sync"
"testing"
"time"
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
"github.com/stretchr/testify/require"
)
const testSendReplyChannel channel.ChannelType = "send_reply_test"
var (
testSendReplyProviderInstance = &testSendReplyProvider{}
testSendReplyRegisterOnce sync.Once
)
type testSendReplyProvider struct {
mu sync.Mutex
count int
err error
externalID string
messageID uint
contactID uint
}
func registerTestSendReplyProvider() *testSendReplyProvider {
testSendReplyRegisterOnce.Do(func() {
_ = channel.Register(testSendReplyProviderInstance)
})
testSendReplyProviderInstance.reset()
return testSendReplyProviderInstance
}
func (p *testSendReplyProvider) reset() {
p.mu.Lock()
defer p.mu.Unlock()
p.count = 0
p.err = nil
p.externalID = "external-message-1"
p.messageID = 0
p.contactID = 0
}
func (p *testSendReplyProvider) Type() channel.ChannelType { return testSendReplyChannel }
func (p *testSendReplyProvider) Name() string { return "Send Reply Test" }
func (p *testSendReplyProvider) Description() string { return "test provider" }
func (p *testSendReplyProvider) ConfigSchema() *channel.ConfigSchemaDefinition {
return &channel.ConfigSchemaDefinition{Type: "object"}
}
func (p *testSendReplyProvider) ValidateConfig(ctx context.Context, config channel.ChannelConfig) error {
return nil
}
func (p *testSendReplyProvider) DefaultConfig() channel.ChannelConfig { return channel.ChannelConfig{} }
func (p *testSendReplyProvider) OnCreate(ctx context.Context, inbox *model.Inbox, config channel.ChannelConfig) (channel.ChannelConfig, error) {
return config, nil
}
func (p *testSendReplyProvider) OnDestroy(ctx context.Context, inbox *model.Inbox, config channel.ChannelConfig) error {
return nil
}
func (p *testSendReplyProvider) ProcessIncoming(ctx context.Context, inbox *model.Inbox, rawPayload []byte) (*channel.IncomingMessage, error) {
return nil, nil
}
func (p *testSendReplyProvider) ValidateWebhookRequest(ctx context.Context, inbox *model.Inbox, request *channel.WebhookRequest) error {
return nil
}
func (p *testSendReplyProvider) SendMessage(ctx context.Context, inbox *model.Inbox, message *model.Message, contact *model.Contact) (*channel.SendResult, error) {
p.mu.Lock()
defer p.mu.Unlock()
p.count++
p.messageID = message.ID
p.contactID = contact.ID
if p.err != nil {
return nil, p.err
}
return &channel.SendResult{ExternalID: p.externalID, DeliveredAt: time.Now()}, nil
}
func (p *testSendReplyProvider) GetContactProfile(ctx context.Context, inbox *model.Inbox, contactSource string) (*channel.ContactProfile, error) {
return nil, nil
}
func (p *testSendReplyProvider) Capabilities() channel.ChannelCapabilities {
return channel.ChannelCapabilities{SupportsDeliveryStatus: true}
}
func TestMessageDeliveryWorker_CreateOutgoingQueuesSendReply(t *testing.T) {
provider := registerTestSendReplyProvider()
db, repo, _, svc := setupMessageServiceWithDefaultLLM(t)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 20, 0, 0, 0, time.UTC) }))
svc.SetWorkerPool(wp)
account := createTestAccount(t, db)
user := createTestUser(t, db, account.ID)
inbox := createTestInbox(t, db, account.ID, string(testSendReplyChannel))
contact := createTestContact(t, db, account.ID)
conversation := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
message, err := svc.Create(context.Background(), account.ID, user.ID, CreateMessageRequest{
ConversationID: conversation.ID,
Content: "hello from agent",
MessageType: "outgoing",
})
require.NoError(t, err)
var queued model.BackgroundJob
require.NoError(t, db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeMessageSendReply, "high", model.BackgroundJobStatusQueued).First(&queued).Error)
processed, err := wp.ProcessOne(context.Background())
require.NoError(t, err)
require.True(t, processed)
updated, err := repo.FindByID(context.Background(), message.ID)
require.NoError(t, err)
require.Equal(t, "external-message-1", updated.SourceID)
require.Equal(t, string(model.MessageStatusSent), updated.Status)
provider.mu.Lock()
require.Equal(t, 1, provider.count)
require.Equal(t, message.ID, provider.messageID)
require.Equal(t, contact.ID, provider.contactID)
provider.mu.Unlock()
}
func TestMessageDeliveryWorker_SkipsMessagesAlreadySentToProvider(t *testing.T) {
provider := registerTestSendReplyProvider()
db := setupServiceTestDB(t)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 20, 15, 0, 0, time.UTC) }))
RegisterMessageDeliveryJobs(wp, db, nil)
account := createTestAccount(t, db)
inbox := createTestInbox(t, db, account.ID, string(testSendReplyChannel))
contact := createTestContact(t, db, account.ID)
conversation := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
message := createTestMessage(t, db, account.ID, inbox.ID, conversation.ID, func(m *model.Message) {
m.MessageType = string(model.MessageTypeOutgoing)
m.SourceID = "existing-external-id"
})
_, err := EnqueueSendReply(context.Background(), wp, message.ID)
require.NoError(t, err)
processed, err := wp.ProcessOne(context.Background())
require.NoError(t, err)
require.True(t, processed)
provider.mu.Lock()
require.Equal(t, 0, provider.count)
provider.mu.Unlock()
}
func TestMessageDeliveryWorker_RetriesProviderFailuresAndMarksMessageFailed(t *testing.T) {
provider := registerTestSendReplyProvider()
provider.mu.Lock()
provider.err = errors.New("provider rejected message")
provider.mu.Unlock()
db := setupServiceTestDB(t)
wp := worker.NewWorkerPoolWithOptions(db,
worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 20, 30, 0, 0, time.UTC) }),
worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }),
)
RegisterMessageDeliveryJobs(wp, db, nil)
account := createTestAccount(t, db)
inbox := createTestInbox(t, db, account.ID, string(testSendReplyChannel))
contact := createTestContact(t, db, account.ID)
conversation := createTestConversation(t, db, account.ID, inbox.ID, contact.ID)
message := createTestMessage(t, db, account.ID, inbox.ID, conversation.ID, func(m *model.Message) {
m.MessageType = string(model.MessageTypeOutgoing)
m.Status = string(model.MessageStatusSent)
})
_, err := EnqueueSendReply(context.Background(), wp, message.ID)
require.NoError(t, err)
processed, err := wp.ProcessOne(context.Background())
require.Error(t, err)
require.True(t, processed)
var job model.BackgroundJob
require.NoError(t, db.Where("job_type = ?", TaskTypeMessageSendReply).First(&job).Error)
require.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
var updated model.Message
require.NoError(t, db.First(&updated, message.ID).Error)
require.Equal(t, string(model.MessageStatusFailed), updated.Status)
attrs := map[string]any{}
require.NoError(t, json.Unmarshal(updated.ContentAttributes, &attrs))
require.Equal(t, "provider rejected message", attrs["external_error"])
}
func TestMessageDeliveryWorker_LoadsMissingMessagesAsRetryableFailures(t *testing.T) {
db := setupServiceTestDB(t)
wp := worker.NewWorkerPoolWithOptions(db,
worker.WithNow(func() time.Time { return time.Date(2026, 6, 5, 20, 45, 0, 0, time.UTC) }),
worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }),
)
RegisterMessageDeliveryJobs(wp, db, nil)
_, err := EnqueueSendReply(context.Background(), wp, 9999)
require.NoError(t, err)
processed, err := wp.ProcessOne(context.Background())
require.Error(t, err)
require.True(t, processed)
var job model.BackgroundJob
require.NoError(t, db.Where("job_type = ?", TaskTypeMessageSendReply).First(&job).Error)
require.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
}