feat(inboxes): queue template sync jobs
This commit is contained in:
@@ -510,6 +510,8 @@ func Bootstrap(env string) (*App, error) {
|
||||
conversationParticipantService := service.NewConversationParticipantService(conversationParticipantRepo, conversationRepo)
|
||||
draftMessageService := service.NewDraftMessageService(draftMessageRepo, conversationRepo)
|
||||
inboxService := service.NewInboxService(inboxRepo, agentBotInboxRepo, agentBotRepo, campaignRepo, webhookSubRepo, waService, waRepo)
|
||||
inboxService.SetWorkerPool(workerPool)
|
||||
service.RegisterInboxTemplateSyncJobs(workerPool, inboxService)
|
||||
igRepo := repository.NewChannelInstagramRepo(db)
|
||||
igService := service.NewChannelInstagramService(igRepo, igProvider)
|
||||
fbChannelRepo := repository.NewChannelFacebookRepo(db)
|
||||
|
||||
@@ -656,13 +656,17 @@ func (h *InboxHandler) SyncTemplates(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
templates, svcErr := h.svc.SyncTemplates(c.Request.Context(), accountID, inboxID)
|
||||
svcErr := h.svc.SyncTemplates(c.Request.Context(), accountID, inboxID)
|
||||
if svcErr != nil {
|
||||
if errors.Is(svcErr, service.ErrInboxTemplateSyncWhatsAppOnly) {
|
||||
c.JSON(http.StatusUnprocessableEntity, gin.H{"error": service.InboxTemplateSyncWhatsAppOnlyMessage})
|
||||
return
|
||||
}
|
||||
handleServiceError(c, svcErr)
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, gin.H{"templates": templates})
|
||||
c.JSON(http.StatusOK, gin.H{"message": service.InboxTemplateSyncInitiatedMessage})
|
||||
}
|
||||
|
||||
// RegisterWebhook registers a webhook URL with the channel provider for an inbox.
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package v1
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
@@ -8,8 +9,17 @@ import (
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
|
||||
whatsappchannel "github.com/gochat/gochat/internal/channel/whatsapp"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
channelmodel "github.com/gochat/gochat/internal/model/channel"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/service"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
)
|
||||
|
||||
// setupInboxMemberActionRouter creates a test router with inbox member-action routes.
|
||||
@@ -123,6 +133,39 @@ func TestInboxSyncTemplates_BadInboxID(t *testing.T) {
|
||||
assert.Contains(t, parseJSONError(w.Body.Bytes()), "invalid inbox id")
|
||||
}
|
||||
|
||||
func TestInboxSyncTemplates_ChatwootQueuedResponse(t *testing.T) {
|
||||
db, err := gorm.Open(sqlite.Open("file:inbox_sync_templates_handler?mode=memory&cache=shared"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
sqlDB, dbErr := db.DB()
|
||||
if dbErr == nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
})
|
||||
require.NoError(t, db.AutoMigrate(&model.Account{}, &model.Inbox{}, &channelmodel.ChannelWhatsApp{}, &model.BackgroundJob{}))
|
||||
account := &model.Account{Name: "Sync Templates", Locale: "en", Active: true}
|
||||
require.NoError(t, db.Create(account).Error)
|
||||
inbox := &model.Inbox{AccountID: account.ID, Name: "WhatsApp", ChannelType: "whatsapp", ChannelID: 1}
|
||||
require.NoError(t, db.Create(inbox).Error)
|
||||
channel := &channelmodel.ChannelWhatsApp{AccountID: account.ID, InboxID: inbox.ID, PhoneNumber: "+1555010000", PhoneNumberID: "phone-1", BusinessAccountID: "waba-1", AccessToken: "token", Provider: "whatsapp_cloud"}
|
||||
require.NoError(t, db.Create(channel).Error)
|
||||
|
||||
wp := worker.NewWorkerPool(db)
|
||||
svc := service.NewInboxService(repository.NewInboxRepo(db), nil, nil, nil, nil, nil, whatsappchannel.NewRepository(db))
|
||||
svc.SetWorkerPool(wp)
|
||||
handler := NewInboxHandler(svc)
|
||||
router := setupInboxMemberActionRouter(handler)
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
req, _ := http.NewRequest("POST", "/api/v1/accounts/1/inboxes/1/sync_templates", bytes.NewReader(nil))
|
||||
router.ServeHTTP(w, req)
|
||||
require.Equal(t, http.StatusOK, w.Code, w.Body.String())
|
||||
var body map[string]any
|
||||
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &body))
|
||||
assert.Equal(t, service.InboxTemplateSyncInitiatedMessage, body["message"])
|
||||
assert.NotContains(t, body, "templates")
|
||||
}
|
||||
|
||||
// ========================================
|
||||
// RegisterWebhook — param validation tests
|
||||
// ========================================
|
||||
|
||||
@@ -15,15 +15,20 @@ import (
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
channelmodel "github.com/gochat/gochat/internal/model/channel"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
pkgvalidator "github.com/gochat/gochat/pkg/validator"
|
||||
)
|
||||
|
||||
const InboxLimitExceededMessage = "Account limit exceeded. Upgrade to a higher plan"
|
||||
const InboxHealthWhatsAppCloudOnlyMessage = "Health data only available for WhatsApp Cloud API channels"
|
||||
const InboxTemplateSyncInitiatedMessage = "Template sync initiated successfully"
|
||||
const InboxTemplateSyncWhatsAppOnlyMessage = "Template sync is only available for WhatsApp channels"
|
||||
const TaskTypeInboxSyncTemplates = "inbox:sync_templates"
|
||||
|
||||
var ErrInboxLimitExceeded = errors.New(InboxLimitExceededMessage)
|
||||
var ErrInboxHealthWhatsAppCloudOnly = errors.New(InboxHealthWhatsAppCloudOnlyMessage)
|
||||
var ErrInboxTemplateSyncWhatsAppOnly = errors.New(InboxTemplateSyncWhatsAppOnlyMessage)
|
||||
|
||||
type WhatsAppChannelService interface {
|
||||
FetchMessageTemplates(ctx context.Context, channel *channelmodel.ChannelWhatsApp) ([]interface{}, error)
|
||||
@@ -41,6 +46,7 @@ type InboxService struct {
|
||||
webhookSubRepo *repository.WebhookSubscriptionRepo
|
||||
whatsappService WhatsAppChannelService
|
||||
whatsappRepo *whatsapp.Repository
|
||||
worker *worker.WorkerPool
|
||||
}
|
||||
|
||||
// NewInboxService creates a new Inbox service.
|
||||
@@ -69,6 +75,10 @@ func (s *InboxService) Ready() bool {
|
||||
return s != nil && s.repo != nil
|
||||
}
|
||||
|
||||
func (s *InboxService) SetWorkerPool(wp *worker.WorkerPool) {
|
||||
s.worker = wp
|
||||
}
|
||||
|
||||
// ListByAccount retrieves all inboxes for an account.
|
||||
func (s *InboxService) ListByAccount(ctx context.Context, accountID uint, offset, limit int) ([]model.Inbox, int64, error) {
|
||||
return s.repo.FindByAccount(ctx, accountID, offset, limit)
|
||||
@@ -1625,36 +1635,28 @@ func (s *InboxService) Health(ctx context.Context, accountID, inboxID uint) (map
|
||||
return s.fetchWhatsAppHealthStatus(ctx, waChannel)
|
||||
}
|
||||
|
||||
// SyncTemplates syncs message templates for an inbox's channel (currently WhatsApp only).
|
||||
// For WhatsApp, this calls the WhatsApp Business API to fetch available templates
|
||||
// and stores them in the channel's message_templates field.
|
||||
// SyncTemplates queues message template sync for an inbox's WhatsApp channel.
|
||||
// Reference: Chatwoot InboxesController#sync_templates (POST member action, WhatsApp only)
|
||||
func (s *InboxService) SyncTemplates(ctx context.Context, accountID, inboxID uint) ([]interface{}, error) {
|
||||
func (s *InboxService) SyncTemplates(ctx context.Context, accountID, inboxID uint) error {
|
||||
inbox, err := s.repo.FindByAccountAndID(ctx, accountID, inboxID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("inbox not found: %w", err)
|
||||
return fmt.Errorf("inbox not found: %w", err)
|
||||
}
|
||||
|
||||
if inbox.ChannelType != "whatsapp" {
|
||||
return nil, fmt.Errorf("sync_templates is only supported for WhatsApp inboxes")
|
||||
return ErrInboxTemplateSyncWhatsAppOnly
|
||||
}
|
||||
|
||||
// Get the WhatsApp channel record
|
||||
waChannel, err := s.getWhatsAppChannel(ctx, inbox.ID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to get WhatsApp channel: %w", err)
|
||||
return fmt.Errorf("failed to get WhatsApp channel: %w", err)
|
||||
}
|
||||
if s.worker == nil {
|
||||
return fmt.Errorf("worker pool not available")
|
||||
}
|
||||
|
||||
// Fetch templates from the WhatsApp Business API
|
||||
templates, err := s.fetchWhatsAppTemplates(ctx, waChannel)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to fetch WhatsApp templates: %w", err)
|
||||
}
|
||||
|
||||
applogger.L().Infof("Synced WhatsApp templates for inbox %d (account_id=%d), got %d templates",
|
||||
inboxID, accountID, len(templates))
|
||||
|
||||
return templates, nil
|
||||
_, err = s.worker.Enqueue(ctx, TaskTypeInboxSyncTemplates, inboxTemplateSyncJob{AccountID: accountID, InboxID: inbox.ID, ChannelID: waChannel.ID}, worker.WithQueue("low"), worker.WithMaxAttempts(3))
|
||||
return err
|
||||
}
|
||||
|
||||
// RegisterWebhookRequest represents the request body for registering a webhook on an inbox.
|
||||
@@ -1752,6 +1754,51 @@ func (s *InboxService) fetchWhatsAppTemplates(ctx context.Context, waChannel *ch
|
||||
return s.whatsappService.FetchMessageTemplates(ctx, waChannel)
|
||||
}
|
||||
|
||||
type inboxTemplateSyncJob struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
InboxID uint `json:"inbox_id"`
|
||||
ChannelID uint `json:"channel_id"`
|
||||
}
|
||||
|
||||
func RegisterInboxTemplateSyncJobs(wp *worker.WorkerPool, svc *InboxService) {
|
||||
if wp == nil || svc == nil {
|
||||
return
|
||||
}
|
||||
wp.Register(TaskTypeInboxSyncTemplates, svc.performTemplateSyncJob)
|
||||
}
|
||||
|
||||
func (s *InboxService) performTemplateSyncJob(ctx context.Context, job *model.BackgroundJob) error {
|
||||
var payload inboxTemplateSyncJob
|
||||
if len(job.Payload) > 0 {
|
||||
if err := json.Unmarshal(job.Payload, &payload); err != nil {
|
||||
return fmt.Errorf("unmarshal inbox template sync job: %w", err)
|
||||
}
|
||||
}
|
||||
if payload.AccountID == 0 || payload.InboxID == 0 {
|
||||
return fmt.Errorf("invalid inbox template sync job payload: %#v", payload)
|
||||
}
|
||||
inbox, err := s.repo.FindByAccountAndID(ctx, payload.AccountID, payload.InboxID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("inbox not found: %w", err)
|
||||
}
|
||||
if inbox.ChannelType != "whatsapp" {
|
||||
return ErrInboxTemplateSyncWhatsAppOnly
|
||||
}
|
||||
waChannel, err := s.getWhatsAppChannel(ctx, inbox.ID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to get WhatsApp channel: %w", err)
|
||||
}
|
||||
if payload.ChannelID != 0 && payload.ChannelID != waChannel.ID {
|
||||
return fmt.Errorf("WhatsApp channel mismatch for inbox_id=%d", inbox.ID)
|
||||
}
|
||||
templates, err := s.fetchWhatsAppTemplates(ctx, waChannel)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to fetch WhatsApp templates: %w", err)
|
||||
}
|
||||
applogger.L().Infof("Synced WhatsApp templates for inbox %d (account_id=%d), got %d templates", inbox.ID, inbox.AccountID, len(templates))
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *InboxService) fetchWhatsAppHealthStatus(ctx context.Context, waChannel *channelmodel.ChannelWhatsApp) (map[string]interface{}, error) {
|
||||
if s.whatsappService == nil {
|
||||
return nil, fmt.Errorf("WhatsApp service not available")
|
||||
|
||||
@@ -2,7 +2,7 @@ package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
channelmodel "github.com/gochat/gochat/internal/model/channel"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
"github.com/gochat/gochat/internal/worker"
|
||||
)
|
||||
|
||||
// ========== Test Setup ==========
|
||||
@@ -37,6 +38,7 @@ func setupInboxServiceTest(t *testing.T) (*InboxService, *gorm.DB) {
|
||||
&model.Inbox{},
|
||||
&model.AgentBotInbox{},
|
||||
&model.WebhookSubscription{},
|
||||
&model.BackgroundJob{},
|
||||
&channelmodel.ChannelWhatsApp{},
|
||||
), "failed to auto-migrate")
|
||||
|
||||
@@ -58,12 +60,19 @@ func setupInboxServiceTest(t *testing.T) (*InboxService, *gorm.DB) {
|
||||
type fakeInboxWhatsAppService struct {
|
||||
healthPayload map[string]interface{}
|
||||
healthErr error
|
||||
templates []interface{}
|
||||
templateErr error
|
||||
fetchCalls int
|
||||
webhookURL string
|
||||
webhookErr error
|
||||
}
|
||||
|
||||
func (f *fakeInboxWhatsAppService) FetchMessageTemplates(context.Context, *channelmodel.ChannelWhatsApp) ([]interface{}, error) {
|
||||
return nil, errors.New("not used")
|
||||
f.fetchCalls++
|
||||
if f.templateErr != nil {
|
||||
return nil, f.templateErr
|
||||
}
|
||||
return f.templates, nil
|
||||
}
|
||||
|
||||
func (f *fakeInboxWhatsAppService) FetchHealthStatus(_ context.Context, _ *channelmodel.ChannelWhatsApp) (map[string]interface{}, error) {
|
||||
@@ -343,34 +352,52 @@ func TestInboxService_SyncTemplates_NonWhatsAppInbox(t *testing.T) {
|
||||
svc, db := setupInboxServiceTest(t)
|
||||
account, inbox := createInboxTestPrereqs(t, db, "api")
|
||||
|
||||
templates, err := svc.SyncTemplates(context.Background(), account.ID, inbox.ID)
|
||||
err := svc.SyncTemplates(context.Background(), account.ID, inbox.ID)
|
||||
assert.Error(t, err)
|
||||
assert.Nil(t, templates)
|
||||
assert.Contains(t, err.Error(), "only supported for WhatsApp")
|
||||
assert.Contains(t, err.Error(), "Template sync is only available for WhatsApp channels")
|
||||
}
|
||||
|
||||
func TestInboxService_SyncTemplates_InboxNotFound(t *testing.T) {
|
||||
svc, _ := setupInboxServiceTest(t)
|
||||
|
||||
templates, err := svc.SyncTemplates(context.Background(), 9999, 9999)
|
||||
err := svc.SyncTemplates(context.Background(), 9999, 9999)
|
||||
assert.Error(t, err)
|
||||
assert.Nil(t, templates)
|
||||
assert.Contains(t, err.Error(), "inbox not found")
|
||||
}
|
||||
|
||||
// SyncTemplates for WhatsApp inboxes requires WhatsApp service/repo which are nil in tests.
|
||||
// The WhatsApp path will fail at getWhatsAppChannel — we verify that it correctly
|
||||
// rejects non-WhatsApp inboxes first (tested above), and for WhatsApp inboxes without
|
||||
// the WhatsApp service, it should error out at the channel lookup step.
|
||||
|
||||
func TestInboxService_SyncTemplates_WhatsAppInbox_NoWaService(t *testing.T) {
|
||||
func TestInboxService_SyncTemplates_WhatsAppQueuesWorkerJob(t *testing.T) {
|
||||
svc, db := setupInboxServiceTest(t)
|
||||
account, inbox := createInboxTestPrereqs(t, db, "whatsapp")
|
||||
account, inbox, channel := createWhatsAppInboxTestPrereqs(t, db, "whatsapp_cloud")
|
||||
wp := worker.NewWorkerPool(db)
|
||||
svc.SetWorkerPool(wp)
|
||||
svc.whatsappService = &fakeInboxWhatsAppService{}
|
||||
|
||||
templates, err := svc.SyncTemplates(context.Background(), account.ID, inbox.ID)
|
||||
// Will fail because whatsapp service/repo are nil
|
||||
assert.Error(t, err)
|
||||
assert.Nil(t, templates)
|
||||
err := svc.SyncTemplates(context.Background(), account.ID, inbox.ID)
|
||||
require.NoError(t, err)
|
||||
|
||||
var job model.BackgroundJob
|
||||
require.NoError(t, db.Where("job_type = ? AND queue = ?", TaskTypeInboxSyncTemplates, "low").First(&job).Error)
|
||||
var payload inboxTemplateSyncJob
|
||||
require.NoError(t, json.Unmarshal(job.Payload, &payload))
|
||||
assert.Equal(t, account.ID, payload.AccountID)
|
||||
assert.Equal(t, inbox.ID, payload.InboxID)
|
||||
assert.Equal(t, channel.ID, payload.ChannelID)
|
||||
}
|
||||
|
||||
func TestInboxService_SyncTemplatesWorkerFetchesTemplates(t *testing.T) {
|
||||
svc, db := setupInboxServiceTest(t)
|
||||
account, inbox, _ := createWhatsAppInboxTestPrereqs(t, db, "whatsapp_cloud")
|
||||
fake := &fakeInboxWhatsAppService{templates: []interface{}{map[string]interface{}{"name": "hello_world"}}}
|
||||
svc.whatsappService = fake
|
||||
wp := worker.NewWorkerPool(db)
|
||||
svc.SetWorkerPool(wp)
|
||||
RegisterInboxTemplateSyncJobs(wp, svc)
|
||||
|
||||
require.NoError(t, svc.SyncTemplates(context.Background(), account.ID, inbox.ID))
|
||||
processed, err := wp.ProcessOne(context.Background())
|
||||
require.NoError(t, err)
|
||||
assert.True(t, processed)
|
||||
assert.Equal(t, 1, fake.fetchCalls)
|
||||
}
|
||||
|
||||
// ========================================
|
||||
|
||||
Reference in New Issue
Block a user