213 lines
7.8 KiB
Go
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)
|
|
}
|