feat(campaigns): derive chatwoot scheduling
This commit is contained in:
@@ -141,6 +141,12 @@ func (s *CampaignHandlerTestSuite) accountURL() string {
|
||||
return "/api/v1/accounts/" + strconv.FormatUint(uint64(s.account.ID), 10) + "/campaigns"
|
||||
}
|
||||
|
||||
func (s *CampaignHandlerTestSuite) seedInbox(channelType string) *model.Inbox {
|
||||
inbox := &model.Inbox{AccountID: s.account.ID, Name: "Campaign " + channelType + " Inbox", ChannelType: channelType, ChannelID: 1}
|
||||
s.Require().NoError(s.db.Create(inbox).Error)
|
||||
return inbox
|
||||
}
|
||||
|
||||
// Helper: seed a campaign directly into the DB for Get/List/Delete/Update tests
|
||||
func (s *CampaignHandlerTestSuite) seedCampaign(title, message, campaignType string) *campaign.Campaign {
|
||||
c := &campaign.Campaign{
|
||||
@@ -256,12 +262,13 @@ func (s *CampaignHandlerTestSuite) TestGet_Unauthorized() {
|
||||
// ========== Create Tests ==========
|
||||
|
||||
func (s *CampaignHandlerTestSuite) TestCreate_Success() {
|
||||
smsInbox := s.seedInbox("Channel::Sms")
|
||||
body := map[string]interface{}{
|
||||
"inbox_id": s.inbox.ID,
|
||||
"inbox_id": smsInbox.ID,
|
||||
"title": "New Campaign",
|
||||
"message": "Hello from campaign",
|
||||
"campaign_type": "one_off",
|
||||
"enabled": true,
|
||||
"scheduled_at": "2026-06-07T10:30:00Z",
|
||||
"audience": []map[string]interface{}{{"type": "Label", "id": 1}},
|
||||
"trigger_rules": map[string]interface{}{"url": "https://example.com"},
|
||||
"template_params": map[string]interface{}{"name": "value"},
|
||||
@@ -281,11 +288,40 @@ func (s *CampaignHandlerTestSuite) TestCreate_Success() {
|
||||
s.NotContains(resp, "data")
|
||||
s.Equal("New Campaign", resp["title"])
|
||||
s.Equal(float64(s.account.ID), resp["account_id"])
|
||||
s.Equal("one_off", resp["campaign_type"])
|
||||
s.Equal(float64(1780828200), resp["scheduled_at"])
|
||||
s.Greater(resp["id"].(float64), float64(0))
|
||||
s.NotNil(resp["inbox"])
|
||||
s.Contains(resp, "audience")
|
||||
s.Contains(resp, "template_params")
|
||||
}
|
||||
|
||||
func (s *CampaignHandlerTestSuite) TestCreate_LiveChatDefaultsOngoingWithoutCampaignType() {
|
||||
body := map[string]interface{}{
|
||||
"inbox_id": s.inbox.ID,
|
||||
"title": "Live Chat Campaign",
|
||||
"message": "Hello from live chat",
|
||||
"enabled": true,
|
||||
"scheduled_at": "2026-06-07T10:30:00Z",
|
||||
"trigger_rules": map[string]interface{}{"url": "https://example.com", "time_on_page": 10},
|
||||
}
|
||||
bodyBytes := marshalNested("campaign", body)
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
req, _ := http.NewRequest("POST", s.accountURL(), bytes.NewReader(bodyBytes))
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
s.router.ServeHTTP(w, req)
|
||||
|
||||
s.Equal(http.StatusOK, w.Code)
|
||||
|
||||
var resp map[string]interface{}
|
||||
s.Require().NoError(json.Unmarshal(w.Body.Bytes(), &resp))
|
||||
s.Equal("Live Chat Campaign", resp["title"])
|
||||
s.Equal("ongoing", resp["campaign_type"])
|
||||
s.NotContains(resp, "scheduled_at")
|
||||
s.NotContains(resp, "audience")
|
||||
}
|
||||
|
||||
func (s *CampaignHandlerTestSuite) TestCreate_ValidationError() {
|
||||
// Missing required fields (title, message, inbox_id)
|
||||
body := map[string]interface{}{
|
||||
@@ -335,10 +371,13 @@ func (s *CampaignHandlerTestSuite) TestCreate_Unauthorized() {
|
||||
|
||||
func (s *CampaignHandlerTestSuite) TestUpdate_Success() {
|
||||
c := s.seedCampaign("Original Title", "Original message", "ongoing")
|
||||
smsInbox := s.seedInbox("Channel::Sms")
|
||||
|
||||
body := map[string]interface{}{
|
||||
"title": "Updated Title",
|
||||
"message": "Updated message",
|
||||
"title": "Updated Title",
|
||||
"message": "Updated message",
|
||||
"inbox_id": smsInbox.ID,
|
||||
"scheduled_at": "2026-06-07T10:30:00Z",
|
||||
}
|
||||
bodyBytes := marshalNested("campaign", body)
|
||||
|
||||
@@ -353,6 +392,14 @@ func (s *CampaignHandlerTestSuite) TestUpdate_Success() {
|
||||
s.Require().NoError(json.Unmarshal(w.Body.Bytes(), &resp))
|
||||
s.NotContains(resp, "success")
|
||||
s.Equal("Updated Title", resp["title"])
|
||||
s.Equal("one_off", resp["campaign_type"])
|
||||
s.Equal(float64(1780828200), resp["scheduled_at"])
|
||||
|
||||
var updated campaign.Campaign
|
||||
s.Require().NoError(s.db.First(&updated, c.ID).Error)
|
||||
s.Equal(smsInbox.ID, updated.InboxID)
|
||||
s.Require().NotNil(updated.ScheduledAt)
|
||||
s.Equal(int64(1780828200), updated.ScheduledAt.Unix())
|
||||
}
|
||||
|
||||
func (s *CampaignHandlerTestSuite) TestUpdate_NotFound() {
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"gorm.io/gorm"
|
||||
|
||||
"github.com/gochat/gochat/internal/campaign"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
)
|
||||
|
||||
// CampaignRepo implements GORM repository for Campaign.
|
||||
@@ -46,6 +47,26 @@ func (r *CampaignRepo) FindByIDAndAccount(ctx context.Context, id, accountID uin
|
||||
return r.FindByDisplayIDAndAccountOrID(ctx, id, accountID)
|
||||
}
|
||||
|
||||
// FindInboxByIDAndAccount retrieves the campaign inbox scoped to the same account.
|
||||
func (r *CampaignRepo) FindInboxByIDAndAccount(ctx context.Context, inboxID, accountID uint) (*model.Inbox, error) {
|
||||
var inbox model.Inbox
|
||||
if err := r.db.WithContext(ctx).Where("id = ? AND account_id = ?", inboxID, accountID).First(&inbox).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &inbox, nil
|
||||
}
|
||||
|
||||
// AccountHasUser returns whether the user belongs to the account.
|
||||
func (r *CampaignRepo) AccountHasUser(ctx context.Context, accountID, userID uint) (bool, error) {
|
||||
var count int64
|
||||
if err := r.db.WithContext(ctx).Model(&model.AccountUser{}).
|
||||
Where("account_id = ? AND user_id = ?", accountID, userID).
|
||||
Count(&count).Error; err != nil {
|
||||
return false, err
|
||||
}
|
||||
return count > 0, nil
|
||||
}
|
||||
|
||||
// FindByDisplayIDAndAccountOrID retrieves a campaign by Chatwoot display_id, falling back to primary key for legacy callers/tests.
|
||||
func (r *CampaignRepo) FindByDisplayIDAndAccountOrID(ctx context.Context, routeID, accountID uint) (*campaign.Campaign, error) {
|
||||
var c campaign.Campaign
|
||||
|
||||
@@ -4,9 +4,12 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gochat/gochat/internal/campaign"
|
||||
"github.com/gochat/gochat/internal/model"
|
||||
"github.com/gochat/gochat/internal/repository"
|
||||
applogger "github.com/gochat/gochat/pkg/logger"
|
||||
pkgvalidator "github.com/gochat/gochat/pkg/validator"
|
||||
@@ -31,7 +34,7 @@ type CreateCampaignRequest struct {
|
||||
Title string `json:"title" validate:"required,min=2"`
|
||||
Message string `json:"message" validate:"required"`
|
||||
Description string `json:"description,omitempty"`
|
||||
CampaignType string `json:"campaign_type" validate:"required,oneof=ongoing one_off"`
|
||||
CampaignType string `json:"campaign_type" validate:"omitempty,oneof=ongoing one_off"`
|
||||
Audience json.RawMessage `json:"audience,omitempty"`
|
||||
TriggerRules json.RawMessage `json:"trigger_rules,omitempty"`
|
||||
TemplateParams json.RawMessage `json:"template_params,omitempty"`
|
||||
@@ -42,6 +45,8 @@ type CreateCampaignRequest struct {
|
||||
|
||||
// UpdateCampaignRequest is the DTO for updating a campaign.
|
||||
type UpdateCampaignRequest struct {
|
||||
InboxID uint `json:"inbox_id,omitempty"`
|
||||
SenderID *uint `json:"sender_id,omitempty"`
|
||||
Title string `json:"title,omitempty" validate:"omitempty,min=2"`
|
||||
Message string `json:"message,omitempty"`
|
||||
Description string `json:"description,omitempty"`
|
||||
@@ -52,8 +57,35 @@ type UpdateCampaignRequest struct {
|
||||
ScheduledAt *string `json:"scheduled_at,omitempty"`
|
||||
Enabled *bool `json:"enabled,omitempty"`
|
||||
TriggerOnlyDuringBusinessHours *bool `json:"trigger_only_during_business_hours,omitempty"`
|
||||
inboxIDSet bool
|
||||
senderIDSet bool
|
||||
descriptionSet bool
|
||||
scheduledAtSet bool
|
||||
}
|
||||
|
||||
func (r *UpdateCampaignRequest) UnmarshalJSON(data []byte) error {
|
||||
type alias UpdateCampaignRequest
|
||||
var raw map[string]json.RawMessage
|
||||
if err := json.Unmarshal(data, &raw); err != nil {
|
||||
return err
|
||||
}
|
||||
var decoded alias
|
||||
if err := json.Unmarshal(data, &decoded); err != nil {
|
||||
return err
|
||||
}
|
||||
*r = UpdateCampaignRequest(decoded)
|
||||
_, r.inboxIDSet = raw["inbox_id"]
|
||||
_, r.senderIDSet = raw["sender_id"]
|
||||
_, r.descriptionSet = raw["description"]
|
||||
_, r.scheduledAtSet = raw["scheduled_at"]
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r UpdateCampaignRequest) InboxIDSet() bool { return r.inboxIDSet }
|
||||
func (r UpdateCampaignRequest) SenderIDSet() bool { return r.senderIDSet }
|
||||
func (r UpdateCampaignRequest) DescriptionSet() bool { return r.descriptionSet }
|
||||
func (r UpdateCampaignRequest) ScheduledAtSet() bool { return r.scheduledAtSet }
|
||||
|
||||
// List retrieves all campaigns for an account with pagination.
|
||||
func (s *CampaignService) List(ctx context.Context, accountID uint, offset, limit int) ([]campaign.Campaign, int64, error) {
|
||||
return s.campaignRepo.ListByAccount(ctx, accountID, offset, limit)
|
||||
@@ -84,6 +116,18 @@ func (s *CampaignService) Create(ctx context.Context, accountID uint, req Create
|
||||
triggerDuringBH = *req.TriggerOnlyDuringBusinessHours
|
||||
}
|
||||
|
||||
inbox, err := s.campaignRepo.FindInboxByIDAndAccount(ctx, req.InboxID, accountID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("inbox not found: %w", err)
|
||||
}
|
||||
if err := s.validateCampaignSender(ctx, accountID, req.SenderID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
campaignType, scheduledAt, err := deriveCampaignAttributes(*inbox, req.ScheduledAt, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
c := &campaign.Campaign{
|
||||
AccountID: accountID,
|
||||
InboxID: req.InboxID,
|
||||
@@ -92,18 +136,14 @@ func (s *CampaignService) Create(ctx context.Context, accountID uint, req Create
|
||||
Message: req.Message,
|
||||
Description: req.Description,
|
||||
CampaignStatus: campaign.CampaignStatusActive,
|
||||
CampaignType: campaign.CampaignType(req.CampaignType),
|
||||
CampaignType: campaignType,
|
||||
Audience: rawJSONParamString(req.Audience),
|
||||
TriggerRules: rawJSONParamString(req.TriggerRules),
|
||||
TemplateParams: rawJSONParamString(req.TemplateParams),
|
||||
ScheduledAt: scheduledAt,
|
||||
Enabled: enabled,
|
||||
TriggerOnlyDuringBusinessHours: triggerDuringBH,
|
||||
}
|
||||
|
||||
// Parse scheduled_at if provided
|
||||
if req.ScheduledAt != nil && *req.ScheduledAt != "" {
|
||||
// ScheduledAt parsing deferred — stored as string for now
|
||||
// Full implementation would parse ISO8601 to time.Time
|
||||
Inbox: *inbox,
|
||||
}
|
||||
|
||||
if err := s.campaignSvc.Create(ctx, c); err != nil {
|
||||
@@ -124,15 +164,39 @@ func (s *CampaignService) Update(ctx context.Context, id, accountID uint, req Up
|
||||
return nil, fmt.Errorf("campaign not found: %w", err)
|
||||
}
|
||||
|
||||
inbox := &c.Inbox
|
||||
if req.InboxIDSet() {
|
||||
if req.InboxID == 0 {
|
||||
return nil, fmt.Errorf("invalid inbox_id")
|
||||
}
|
||||
resolvedInbox, err := s.campaignRepo.FindInboxByIDAndAccount(ctx, req.InboxID, accountID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("inbox not found: %w", err)
|
||||
}
|
||||
inbox = resolvedInbox
|
||||
c.InboxID = req.InboxID
|
||||
}
|
||||
if req.SenderIDSet() {
|
||||
if err := s.validateCampaignSender(ctx, accountID, req.SenderID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
c.SenderID = req.SenderID
|
||||
}
|
||||
campaignType, scheduledAt, err := deriveCampaignAttributes(*inbox, req.ScheduledAt, c.ScheduledAt)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
c.CampaignType = campaignType
|
||||
c.ScheduledAt = scheduledAt
|
||||
|
||||
if req.Title != "" {
|
||||
c.Title = req.Title
|
||||
}
|
||||
if req.Message != "" {
|
||||
c.Message = req.Message
|
||||
}
|
||||
c.Description = req.Description
|
||||
if req.CampaignType != "" {
|
||||
c.CampaignType = campaign.CampaignType(req.CampaignType)
|
||||
if req.DescriptionSet() {
|
||||
c.Description = req.Description
|
||||
}
|
||||
if len(req.Audience) > 0 {
|
||||
c.Audience = rawJSONParamString(req.Audience)
|
||||
@@ -151,6 +215,8 @@ func (s *CampaignService) Update(ctx context.Context, id, accountID uint, req Up
|
||||
}
|
||||
|
||||
if err := s.campaignSvc.Update(ctx, c.ID, map[string]interface{}{
|
||||
"inbox_id": c.InboxID,
|
||||
"sender_id": c.SenderID,
|
||||
"title": c.Title,
|
||||
"message": c.Message,
|
||||
"description": c.Description,
|
||||
@@ -158,13 +224,18 @@ func (s *CampaignService) Update(ctx context.Context, id, accountID uint, req Up
|
||||
"audience": c.Audience,
|
||||
"trigger_rules": c.TriggerRules,
|
||||
"template_params": c.TemplateParams,
|
||||
"scheduled_at": c.ScheduledAt,
|
||||
"enabled": c.Enabled,
|
||||
"trigger_only_during_business_hours": c.TriggerOnlyDuringBusinessHours,
|
||||
}); err != nil {
|
||||
applogger.L().Errorf("failed to update campaign: %v", err)
|
||||
return nil, fmt.Errorf("failed to update campaign: %w", err)
|
||||
}
|
||||
return c, nil
|
||||
reloadedID := c.DisplayID
|
||||
if reloadedID == 0 {
|
||||
reloadedID = c.ID
|
||||
}
|
||||
return s.campaignRepo.FindByDisplayIDAndAccountOrID(ctx, reloadedID, accountID)
|
||||
}
|
||||
|
||||
// Delete soft-deletes a campaign scoped to an account.
|
||||
@@ -203,6 +274,80 @@ func rawJSONParamString(raw json.RawMessage) string {
|
||||
return trimmed
|
||||
}
|
||||
|
||||
func (s *CampaignService) validateCampaignSender(ctx context.Context, accountID uint, senderID *uint) error {
|
||||
if senderID == nil || *senderID == 0 {
|
||||
return nil
|
||||
}
|
||||
ok, err := s.campaignRepo.AccountHasUser(ctx, accountID, *senderID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("validate sender: %w", err)
|
||||
}
|
||||
if !ok {
|
||||
return fmt.Errorf("invalid sender_id: must belong to the same account as the campaign")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func deriveCampaignAttributes(inbox model.Inbox, scheduledAt *string, current *time.Time) (campaign.CampaignType, *time.Time, error) {
|
||||
if campaignInboxIsOneOff(inbox.ChannelType) {
|
||||
parsed, err := parseCampaignScheduledAt(scheduledAt)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
if parsed == nil {
|
||||
if current != nil {
|
||||
copy := current.UTC()
|
||||
parsed = ©
|
||||
} else {
|
||||
now := time.Now().UTC()
|
||||
parsed = &now
|
||||
}
|
||||
}
|
||||
return campaign.CampaignTypeOneOff, parsed, nil
|
||||
}
|
||||
if campaignInboxIsOngoing(inbox.ChannelType) {
|
||||
return campaign.CampaignTypeOngoing, nil, nil
|
||||
}
|
||||
return "", nil, fmt.Errorf("invalid inbox: Unsupported Inbox type")
|
||||
}
|
||||
|
||||
func campaignInboxIsOneOff(channelType string) bool {
|
||||
normalized := normalizeCampaignInboxType(channelType)
|
||||
return normalized == "sms" || normalized == "twiliosms" || normalized == "whatsapp"
|
||||
}
|
||||
|
||||
func campaignInboxIsOngoing(channelType string) bool {
|
||||
normalized := normalizeCampaignInboxType(channelType)
|
||||
return normalized == "webwidget" || normalized == "website"
|
||||
}
|
||||
|
||||
func normalizeCampaignInboxType(channelType string) string {
|
||||
normalized := strings.ToLower(strings.TrimSpace(channelType))
|
||||
normalized = strings.TrimPrefix(normalized, "channel::")
|
||||
normalized = strings.ReplaceAll(normalized, "_", "")
|
||||
normalized = strings.ReplaceAll(normalized, " ", "")
|
||||
normalized = strings.ReplaceAll(normalized, "-", "")
|
||||
return normalized
|
||||
}
|
||||
|
||||
func parseCampaignScheduledAt(raw *string) (*time.Time, error) {
|
||||
if raw == nil || strings.TrimSpace(*raw) == "" {
|
||||
return nil, nil
|
||||
}
|
||||
trimmed := strings.TrimSpace(*raw)
|
||||
if unix, err := strconv.ParseInt(trimmed, 10, 64); err == nil {
|
||||
parsed := time.Unix(unix, 0).UTC()
|
||||
return &parsed, nil
|
||||
}
|
||||
for _, layout := range []string{time.RFC3339Nano, time.RFC3339, "2006-01-02 15:04:05 -0700", "2006-01-02 15:04:05"} {
|
||||
if parsed, err := time.Parse(layout, trimmed); err == nil {
|
||||
parsed = parsed.UTC()
|
||||
return &parsed, nil
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("invalid scheduled_at")
|
||||
}
|
||||
|
||||
// Stop marks a campaign as completed.
|
||||
func (s *CampaignService) Stop(ctx context.Context, id, accountID uint) error {
|
||||
c, err := s.campaignRepo.FindByIDAndAccount(ctx, id, accountID)
|
||||
|
||||
Reference in New Issue
Block a user