feat(messages): queue delivery statuses

This commit is contained in:
2026-06-05 19:05:45 +08:00
parent 1a4b4a2d29
commit dec8854ec9
11 changed files with 711 additions and 57 deletions
@@ -39,6 +39,7 @@ import (
fbchannel "github.com/gochat/gochat/internal/channel/facebook"
"github.com/gochat/gochat/internal/model"
channelmodel "github.com/gochat/gochat/internal/model/channel"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"gorm.io/gorm"
)
@@ -68,6 +69,13 @@ func NewFacebookWebhookHandler(
}
}
func (h *FacebookWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *FacebookWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetWorkerPool(wp)
}
return h
}
// HandleFacebookVerification handles the GET webhook verification request from Facebook.
// URL pattern: /webhooks/facebook/:inbox_id
// Method: GET
+83 -6
View File
@@ -12,6 +12,7 @@ import (
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
// IncomingPersister is the durable boundary after provider-specific webhook parsing.
@@ -19,6 +20,7 @@ import (
type IncomingPersister struct {
db *gorm.DB
dispatcher *channel.Dispatcher
worker *worker.WorkerPool
}
type IncomingPersistResult struct {
@@ -31,6 +33,15 @@ type IncomingPersistResult struct {
ConversationCreated bool
}
func (p *IncomingPersister) SetWorkerPool(wp *worker.WorkerPool) *IncomingPersister {
if p == nil || wp == nil {
return p
}
p.worker = wp
p.registerJobs(wp)
return p
}
func NewIncomingPersister(db *gorm.DB, dispatcher ...*channel.Dispatcher) *IncomingPersister {
if db == nil {
return nil
@@ -102,6 +113,21 @@ func (p *IncomingPersister) PersistIncoming(ctx context.Context, inbox *model.In
// UpdateMessageStatus applies provider delivery/read/failed receipts to existing messages.
func (p *IncomingPersister) UpdateMessageStatus(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time) error {
return p.UpdateMessageStatusWithError(ctx, inbox, sourceID, status, occurredAt, "")
}
// UpdateMessageStatusWithError applies provider receipts and records provider failure details.
func (p *IncomingPersister) UpdateMessageStatusWithError(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time, externalError string) error {
if p == nil || p.db == nil || inbox == nil || sourceID == "" {
return nil
}
if p.worker != nil {
return p.enqueueMessageStatusUpdate(ctx, inbox.ID, sourceID, status, occurredAt, externalError)
}
return p.performMessageStatusUpdate(ctx, inbox, sourceID, status, occurredAt, externalError)
}
func (p *IncomingPersister) performMessageStatusUpdate(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time, externalError string) error {
if p == nil || p.db == nil || inbox == nil || sourceID == "" {
return nil
}
@@ -114,15 +140,27 @@ func (p *IncomingPersister) UpdateMessageStatus(ctx context.Context, inbox *mode
}
return err
}
if err := tx.Model(&message).Update("status", string(status)).Error; err != nil {
if !validProviderMessageStatusTransition(model.MessageStatus(message.Status), status) {
return nil
}
updates := map[string]any{
"status": string(status),
"content_attributes": setProviderMessageExternalError(message.ContentAttributes, status, externalError),
}
if err := tx.Model(&message).Updates(updates).Error; err != nil {
return err
}
message.Status = string(status)
if message.SenderID == nil {
message.ContentAttributes = updates["content_attributes"].(datatypes.JSON)
contactID, err := p.deliveryStatusContactID(ctx, tx, &message)
if err != nil {
return err
}
if contactID == 0 {
p.dispatchMessageStatusEvent(ctx, inbox, &message)
return nil
}
if err := p.upsertDeliveryStatus(ctx, tx, &message, *message.SenderID, status, occurredAt); err != nil {
if err := p.upsertDeliveryStatus(ctx, tx, &message, contactID, status, occurredAt); err != nil {
return err
}
p.dispatchMessageStatusEvent(ctx, inbox, &message)
@@ -135,6 +173,16 @@ func (p *IncomingPersister) UpdateContactConversationMessagesStatus(ctx context.
if p == nil || p.db == nil || inbox == nil || contactSourceID == "" {
return nil
}
if p.worker != nil {
return p.enqueueContactMessagesStatusUpdate(ctx, inbox.ID, contactSourceID, status, occurredAt)
}
return p.performContactMessagesStatusUpdate(ctx, inbox, contactSourceID, status, occurredAt)
}
func (p *IncomingPersister) performContactMessagesStatusUpdate(ctx context.Context, inbox *model.Inbox, contactSourceID string, status model.MessageStatus, occurredAt *time.Time) error {
if p == nil || p.db == nil || inbox == nil || contactSourceID == "" {
return nil
}
return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var contactInbox model.ContactInbox
@@ -157,17 +205,46 @@ func (p *IncomingPersister) UpdateContactConversationMessagesStatus(ctx context.
if err := query.Select("messages.*").Find(&messages).Error; err != nil {
return err
}
if err := query.Update("status", string(status)).Error; err != nil {
return err
}
for i := range messages {
if !validProviderMessageStatusTransition(model.MessageStatus(messages[i].Status), status) {
continue
}
if err := tx.Model(&messages[i]).Updates(map[string]any{
"status": string(status),
"content_attributes": setProviderMessageExternalError(messages[i].ContentAttributes, status, ""),
}).Error; err != nil {
return err
}
messages[i].Status = string(status)
if err := p.upsertDeliveryStatus(ctx, tx, &messages[i], contactInbox.ContactID, status, occurredAt); err != nil {
return err
}
p.dispatchMessageStatusEvent(ctx, inbox, &messages[i])
}
return nil
})
}
func (p *IncomingPersister) deliveryStatusContactID(ctx context.Context, tx *gorm.DB, message *model.Message) (uint, error) {
if message == nil {
return 0, nil
}
if message.ConversationID != 0 {
var conversation model.Conversation
if err := tx.WithContext(ctx).Select("id", "contact_id").First(&conversation, message.ConversationID).Error; err != nil {
if err != gorm.ErrRecordNotFound {
return 0, err
}
} else if conversation.ContactID != 0 {
return conversation.ContactID, nil
}
}
if message.SenderID != nil {
return *message.SenderID, nil
}
return 0, nil
}
func (p *IncomingPersister) upsertDeliveryStatus(ctx context.Context, tx *gorm.DB, message *model.Message, contactID uint, status model.MessageStatus, occurredAt *time.Time) error {
var delivery model.DeliveryStatus
err := tx.WithContext(ctx).Where("message_id = ? AND contact_id = ?", message.ID, contactID).First(&delivery).Error
@@ -0,0 +1,194 @@
package webhook
import (
"context"
"encoding/json"
"fmt"
"strings"
"time"
"gorm.io/datatypes"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
const (
TaskTypeProviderMessageStatusUpdate = "webhook:message_status_update"
TaskTypeProviderContactMessagesStatusUpdate = "webhook:contact_messages_status_update"
)
type providerMessageStatusUpdateJob struct {
InboxID uint `json:"inbox_id"`
SourceID string `json:"source_id"`
Status model.MessageStatus `json:"status"`
OccurredAtUnixNano int64 `json:"occurred_at_unix_nano,omitempty"`
ExternalError string `json:"external_error,omitempty"`
}
type providerContactMessagesStatusUpdateJob struct {
InboxID uint `json:"inbox_id"`
ContactSourceID string `json:"contact_source_id"`
Status model.MessageStatus `json:"status"`
OccurredAtUnixNano int64 `json:"occurred_at_unix_nano,omitempty"`
}
func (p *IncomingPersister) registerJobs(wp *worker.WorkerPool) {
if p == nil || wp == nil {
return
}
wp.Register(TaskTypeProviderMessageStatusUpdate, p.performMessageStatusUpdateJob)
wp.Register(TaskTypeProviderContactMessagesStatusUpdate, p.performContactMessagesStatusUpdateJob)
}
func (p *IncomingPersister) enqueueMessageStatusUpdate(ctx context.Context, inboxID uint, sourceID string, status model.MessageStatus, occurredAt *time.Time, externalError string) error {
if p == nil || p.worker == nil || inboxID == 0 || sourceID == "" {
return nil
}
if !validProviderMessageStatus(status) {
return nil
}
payload := providerMessageStatusUpdateJob{
InboxID: inboxID,
SourceID: sourceID,
Status: status,
OccurredAtUnixNano: timeToUnixNano(occurredAt),
ExternalError: strings.TrimSpace(externalError),
}
_, err := p.worker.Enqueue(ctx, TaskTypeProviderMessageStatusUpdate, payload,
worker.WithQueue("low"),
worker.WithMaxAttempts(3),
worker.WithIdempotencyKey(fmt.Sprintf("webhook:message_status:%d:%s:%s:%d", inboxID, sourceID, status, payload.OccurredAtUnixNano)),
)
return err
}
func (p *IncomingPersister) enqueueContactMessagesStatusUpdate(ctx context.Context, inboxID uint, contactSourceID string, status model.MessageStatus, occurredAt *time.Time) error {
if p == nil || p.worker == nil || inboxID == 0 || contactSourceID == "" {
return nil
}
if !validProviderMessageStatus(status) {
return nil
}
payload := providerContactMessagesStatusUpdateJob{
InboxID: inboxID,
ContactSourceID: contactSourceID,
Status: status,
OccurredAtUnixNano: timeToUnixNano(occurredAt),
}
_, err := p.worker.Enqueue(ctx, TaskTypeProviderContactMessagesStatusUpdate, payload,
worker.WithQueue("low"),
worker.WithMaxAttempts(3),
worker.WithIdempotencyKey(fmt.Sprintf("webhook:contact_messages_status:%d:%s:%s:%d", inboxID, contactSourceID, status, payload.OccurredAtUnixNano)),
)
return err
}
func (p *IncomingPersister) performMessageStatusUpdateJob(ctx context.Context, job *model.BackgroundJob) error {
var payload providerMessageStatusUpdateJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal provider message status job: %w", err)
}
if payload.InboxID == 0 || payload.SourceID == "" || !validProviderMessageStatus(payload.Status) {
return fmt.Errorf("invalid provider message status job payload: %#v", payload)
}
inbox, err := p.loadInboxForStatusJob(ctx, payload.InboxID)
if err != nil {
return err
}
if inbox == nil {
return nil
}
return p.performMessageStatusUpdate(ctx, inbox, payload.SourceID, payload.Status, unixNanoToTime(payload.OccurredAtUnixNano), payload.ExternalError)
}
func (p *IncomingPersister) performContactMessagesStatusUpdateJob(ctx context.Context, job *model.BackgroundJob) error {
var payload providerContactMessagesStatusUpdateJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal provider contact messages status job: %w", err)
}
if payload.InboxID == 0 || payload.ContactSourceID == "" || !validProviderMessageStatus(payload.Status) {
return fmt.Errorf("invalid provider contact messages status job payload: %#v", payload)
}
inbox, err := p.loadInboxForStatusJob(ctx, payload.InboxID)
if err != nil {
return err
}
if inbox == nil {
return nil
}
return p.performContactMessagesStatusUpdate(ctx, inbox, payload.ContactSourceID, payload.Status, unixNanoToTime(payload.OccurredAtUnixNano))
}
func (p *IncomingPersister) loadInboxForStatusJob(ctx context.Context, inboxID uint) (*model.Inbox, error) {
var inbox model.Inbox
if err := p.db.WithContext(ctx).First(&inbox, inboxID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil, nil
}
return nil, fmt.Errorf("load status inbox %d: %w", inboxID, err)
}
return &inbox, nil
}
func validProviderMessageStatus(status model.MessageStatus) bool {
switch status {
case model.MessageStatusSent, model.MessageStatusDelivered, model.MessageStatusRead, model.MessageStatusFailed:
return true
default:
return false
}
}
func validProviderMessageStatusTransition(current, next model.MessageStatus) bool {
if !validProviderMessageStatus(next) || current == next {
return validProviderMessageStatus(next)
}
if next == model.MessageStatusFailed || current == model.MessageStatusFailed {
return true
}
return providerMessageStatusRank(next) >= providerMessageStatusRank(current)
}
func providerMessageStatusRank(status model.MessageStatus) int {
switch status {
case model.MessageStatusSent:
return 1
case model.MessageStatusDelivered:
return 2
case model.MessageStatusRead:
return 3
default:
return 0
}
}
func setProviderMessageExternalError(attrs datatypes.JSON, status model.MessageStatus, externalError string) datatypes.JSON {
obj := map[string]any{}
if len(attrs) > 0 {
_ = json.Unmarshal(attrs, &obj)
}
if status == model.MessageStatusFailed && strings.TrimSpace(externalError) != "" {
obj["external_error"] = strings.TrimSpace(externalError)
} else {
delete(obj, "external_error")
}
bytes, _ := json.Marshal(obj)
return datatypes.JSON(bytes)
}
func timeToUnixNano(t *time.Time) int64 {
if t == nil || t.IsZero() {
return 0
}
return t.UTC().UnixNano()
}
func unixNanoToTime(value int64) *time.Time {
if value <= 0 {
return nil
}
t := time.Unix(0, value).UTC()
return &t
}
@@ -22,6 +22,7 @@ import (
tiktokchannel "github.com/gochat/gochat/internal/channel/tiktok"
"github.com/gochat/gochat/internal/model"
channelmodel "github.com/gochat/gochat/internal/model/channel"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"github.com/gin-gonic/gin"
@@ -46,6 +47,13 @@ func NewTikTokWebhookHandler(tiktokWebhook *tiktokchannel.WebhookHandler, pipeli
}
}
func (h *TikTokWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *TikTokWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetWorkerPool(wp)
}
return h
}
// HandleTikTokWebhook processes incoming TikTok webhook HTTP requests.
func (h *TikTokWebhookHandler) HandleTikTokWebhook(c *gin.Context) {
body, err := io.ReadAll(c.Request.Body)
+59 -18
View File
@@ -13,11 +13,13 @@ package webhook
import (
"fmt"
"net/http"
"net/url"
"github.com/gochat/gochat/internal/channel"
twiliochannel "github.com/gochat/gochat/internal/channel/twilio"
"github.com/gochat/gochat/internal/model"
channelmodel "github.com/gochat/gochat/internal/model/channel"
"github.com/gochat/gochat/internal/worker"
applogger "github.com/gochat/gochat/pkg/logger"
"github.com/gin-gonic/gin"
@@ -31,6 +33,13 @@ type TwilioWebhookHandler struct {
persister *IncomingPersister
}
func (h *TwilioWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *TwilioWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetWorkerPool(wp)
}
return h
}
// NewTwilioWebhookHandler creates a Twilio SMS webhook handler for Gin integration.
func NewTwilioWebhookHandler(twilioWebhook *twiliochannel.WebhookHandler, db *gorm.DB, dispatcher ...*channel.Dispatcher) *TwilioWebhookHandler {
return &TwilioWebhookHandler{
@@ -74,33 +83,65 @@ func (h *TwilioWebhookHandler) HandleTwilioInboundSMS(c *gin.Context) {
// HandleTwilioDeliveryStatus processes a Twilio delivery status callback.
func (h *TwilioWebhookHandler) HandleTwilioDeliveryStatus(c *gin.Context) {
phoneNumber := c.Param("phone_number")
if phoneNumber == "" {
applogger.L().Warn("Twilio status webhook: missing phone_number in path")
c.Status(http.StatusOK)
return
}
// Lookup inbox from database
inbox, err := h.lookupInboxByPhoneNumber(phoneNumber)
if err != nil {
applogger.L().Warnf("Twilio status webhook: inbox lookup failed for phone_number %s: %v", phoneNumber, err)
c.Status(http.StatusOK)
return
}
if err := c.Request.ParseForm(); err != nil {
applogger.L().Errorf("Twilio status webhook: parse form failed for inbox %d: %v", inbox.ID, err)
c.Status(http.StatusOK)
applogger.L().Errorf("Twilio status webhook: parse form failed: %v", err)
c.Status(http.StatusNoContent)
return
}
var inbox *model.Inbox
var err error
if phoneNumber != "" {
inbox, err = h.lookupInboxByPhoneNumber(phoneNumber)
} else {
inbox, err = h.lookupDeliveryStatusInbox(c.Request.Form)
}
if err != nil {
applogger.L().Warnf("Twilio status webhook: inbox lookup failed: %v", err)
c.Status(http.StatusNoContent)
return
}
messageSID := c.Request.FormValue("MessageSid")
messageStatus := c.Request.FormValue("MessageStatus")
if mapped, ok := mapTwilioMessageStatus(messageStatus); ok {
if err := h.persister.UpdateMessageStatus(c.Request.Context(), inbox, messageSID, mapped, nil); err != nil {
if err := h.persister.UpdateMessageStatusWithError(c.Request.Context(), inbox, messageSID, mapped, nil, twilioExternalError(c.Request.FormValue("ErrorCode"), c.Request.FormValue("ErrorMessage"), messageStatus)); err != nil {
applogger.L().Errorf("Twilio status webhook: status persistence failed for inbox %d sid=%s status=%s: %v", inbox.ID, messageSID, messageStatus, err)
}
}
c.Status(http.StatusOK)
c.Status(http.StatusNoContent)
}
func twilioExternalError(errorCode, errorMessage, status string) string {
if errorCode == "" || (status != "failed" && status != "undelivered") {
return ""
}
if errorMessage != "" {
return fmt.Sprintf("%s - %s", errorCode, errorMessage)
}
return fmt.Sprintf("Twilio delivery failed with error code %s", errorCode)
}
func (h *TwilioWebhookHandler) lookupDeliveryStatusInbox(params url.Values) (*model.Inbox, error) {
if h.db == nil {
return nil, fmt.Errorf("twilio webhook database is not configured")
}
var twilioChannel channelmodel.ChannelTwilioSMS
query := h.db
if sid := params.Get("MessagingServiceSid"); sid != "" {
query = query.Where(&channelmodel.ChannelTwilioSMS{MessagingServiceSID: sid})
} else if accountSID, from := params.Get("AccountSid"), params.Get("From"); accountSID != "" && from != "" {
query = query.Where(&channelmodel.ChannelTwilioSMS{AccountSID: accountSID, PhoneNumber: from})
} else {
return nil, fmt.Errorf("delivery status missing MessagingServiceSid or AccountSid/From")
}
if err := query.First(&twilioChannel).Error; err != nil {
return nil, err
}
var inbox model.Inbox
if err := h.db.Where("id = ? AND channel_type IN ?", twilioChannel.InboxID, []string{"twilio_sms", "sms"}).First(&inbox).Error; err != nil {
return nil, err
}
return &inbox, nil
}
func mapTwilioMessageStatus(status string) (model.MessageStatus, bool) {
+292 -2
View File
@@ -26,6 +26,7 @@ import (
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/worker"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
@@ -62,6 +63,7 @@ func newWebhookLookupTestDB(t *testing.T) *gorm.DB {
&model.Conversation{},
&model.Message{},
&model.DeliveryStatus{},
&model.BackgroundJob{},
&automation.AutomationRule{},
&automation.AutomationExecution{},
&channelmodel.ChannelTelegram{},
@@ -167,6 +169,193 @@ func TestIncomingPersisterUpdatesMessageStatus(t *testing.T) {
}
}
func TestIncomingPersisterQueuesMessageStatusUpdateWithWorker(t *testing.T) {
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "telegram")
dispatcher := channel.NewDispatcher()
listener := &recordingListener{}
dispatcher.Register(listener)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db, dispatcher).SetWorkerPool(wp)
msg := &channel.IncomingMessage{
ChannelType: channel.ChannelTelegram,
SourceID: "tg-status-job-1",
SenderID: "tg-user-status-job",
SenderName: "Status Worker User",
SenderType: channel.SenderContact,
Content: "status me async",
ContentType: channel.ContentText,
InboxID: inbox.ID,
AccountID: inbox.AccountID,
}
result, err := persister.PersistIncoming(t.Context(), &inbox, msg)
if err != nil {
t.Fatalf("persist incoming: %v", err)
}
occurredAt := time.Now().UTC().Add(-time.Minute)
if err := persister.UpdateMessageStatus(t.Context(), &inbox, "tg-status-job-1", model.MessageStatusDelivered, &occurredAt); err != nil {
t.Fatalf("enqueue status: %v", err)
}
var queued model.BackgroundJob
if err := db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeProviderMessageStatusUpdate, "low", model.BackgroundJobStatusQueued).First(&queued).Error; err != nil {
t.Fatalf("expected queued status job: %v", err)
}
var before model.Message
if err := db.First(&before, result.Message.ID).Error; err != nil {
t.Fatalf("load before message: %v", err)
}
if before.Status != string(model.MessageStatusSent) {
t.Fatalf("expected status to remain sent before worker, got %s", before.Status)
}
processed, err := wp.ProcessOne(t.Context())
if err != nil || !processed {
t.Fatalf("process status job processed=%v err=%v", processed, err)
}
var updated model.Message
if err := db.First(&updated, result.Message.ID).Error; err != nil {
t.Fatalf("load updated message: %v", err)
}
if updated.Status != string(model.MessageStatusDelivered) {
t.Fatalf("expected delivered status, got %s", updated.Status)
}
var delivery model.DeliveryStatus
if err := db.Where("message_id = ? AND contact_id = ?", updated.ID, result.Contact.ID).First(&delivery).Error; err != nil {
t.Fatalf("expected delivery status: %v", err)
}
if delivery.Status != model.MessageStatusDelivered || delivery.DeliveredAt == nil {
t.Fatalf("expected delivered timestamp, got %#v", delivery)
}
if !listenerSaw(listener, channel.EventMessageStatusUpdated) {
t.Fatalf("expected message.status_updated event, got %#v", listener.events)
}
}
func TestIncomingPersisterStatusJobDoesNotDowngradeRead(t *testing.T) {
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "telegram")
contact := model.Contact{AccountID: inbox.AccountID, Name: "Read User"}
if err := db.Create(&contact).Error; err != nil {
t.Fatalf("create contact: %v", err)
}
conversation := model.Conversation{AccountID: inbox.AccountID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType}
if err := db.Create(&conversation).Error; err != nil {
t.Fatalf("create conversation: %v", err)
}
message := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "already read", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusRead), SourceID: "provider-read-1"}
if err := db.Create(&message).Error; err != nil {
t.Fatalf("create message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db).SetWorkerPool(wp)
if err := persister.UpdateMessageStatus(t.Context(), &inbox, "provider-read-1", model.MessageStatusDelivered, nil); err != nil {
t.Fatalf("enqueue downgrade: %v", err)
}
processed, err := wp.ProcessOne(t.Context())
if err != nil || !processed {
t.Fatalf("process downgrade job processed=%v err=%v", processed, err)
}
var updated model.Message
if err := db.First(&updated, message.ID).Error; err != nil {
t.Fatalf("load message: %v", err)
}
if updated.Status != string(model.MessageStatusRead) {
t.Fatalf("expected read to remain read, got %s", updated.Status)
}
}
func TestIncomingPersisterStatusJobStoresExternalError(t *testing.T) {
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "twilio_sms")
contact := model.Contact{AccountID: inbox.AccountID, Name: "SMS User"}
if err := db.Create(&contact).Error; err != nil {
t.Fatalf("create contact: %v", err)
}
conversation := model.Conversation{AccountID: inbox.AccountID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType}
if err := db.Create(&conversation).Error; err != nil {
t.Fatalf("create conversation: %v", err)
}
message := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "failed", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusSent), SourceID: "SMFAIL1"}
if err := db.Create(&message).Error; err != nil {
t.Fatalf("create message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db).SetWorkerPool(wp)
if err := persister.UpdateMessageStatusWithError(t.Context(), &inbox, "SMFAIL1", model.MessageStatusFailed, nil, "30003 - unreachable handset"); err != nil {
t.Fatalf("enqueue failure: %v", err)
}
processed, err := wp.ProcessOne(t.Context())
if err != nil || !processed {
t.Fatalf("process failure job processed=%v err=%v", processed, err)
}
var updated model.Message
if err := db.First(&updated, message.ID).Error; err != nil {
t.Fatalf("load failed message: %v", err)
}
if updated.Status != string(model.MessageStatusFailed) {
t.Fatalf("expected failed status, got %s", updated.Status)
}
attrs := map[string]any{}
if err := json.Unmarshal(updated.ContentAttributes, &attrs); err != nil {
t.Fatalf("unmarshal content attrs: %v", err)
}
if attrs["external_error"] != "30003 - unreachable handset" {
t.Fatalf("expected external_error, got %#v", attrs)
}
}
func TestIncomingPersisterQueuesContactMessagesStatusUpdateWithWorker(t *testing.T) {
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "facebook")
contact := model.Contact{AccountID: inbox.AccountID, Name: "Meta User"}
if err := db.Create(&contact).Error; err != nil {
t.Fatalf("create contact: %v", err)
}
contactInbox := model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: "meta-user-1", PubsubToken: "pub-meta"}
if err := db.Create(&contactInbox).Error; err != nil {
t.Fatalf("create contact inbox: %v", err)
}
conversation := model.Conversation{AccountID: inbox.AccountID, InboxID: inbox.ID, ContactID: contact.ID, ContactInboxID: &contactInbox.ID, Status: string(model.ConversationStatusOpen), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType}
if err := db.Create(&conversation).Error; err != nil {
t.Fatalf("create conversation: %v", err)
}
cutoff := time.Now().UTC()
beforeSent := createWebhookStatusMessage(t, db, inbox, conversation, "sent-before", model.MessageStatusSent, cutoff.Add(-time.Minute))
beforeDelivered := createWebhookStatusMessage(t, db, inbox, conversation, "delivered-before", model.MessageStatusDelivered, cutoff.Add(-30*time.Second))
alreadyRead := createWebhookStatusMessage(t, db, inbox, conversation, "read-before", model.MessageStatusRead, cutoff.Add(-20*time.Second))
afterCutoff := createWebhookStatusMessage(t, db, inbox, conversation, "sent-after", model.MessageStatusSent, cutoff.Add(time.Minute))
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db).SetWorkerPool(wp)
if err := persister.UpdateContactConversationMessagesStatus(t.Context(), &inbox, "meta-user-1", model.MessageStatusRead, &cutoff); err != nil {
t.Fatalf("enqueue contact read: %v", err)
}
var queued model.BackgroundJob
if err := db.Where("job_type = ? AND queue = ?", TaskTypeProviderContactMessagesStatusUpdate, "low").First(&queued).Error; err != nil {
t.Fatalf("expected queued contact status job: %v", err)
}
processed, err := wp.ProcessOne(t.Context())
if err != nil || !processed {
t.Fatalf("process contact status job processed=%v err=%v", processed, err)
}
assertWebhookMessageStatus(t, db, beforeSent.ID, model.MessageStatusRead)
assertWebhookMessageStatus(t, db, beforeDelivered.ID, model.MessageStatusRead)
assertWebhookMessageStatus(t, db, alreadyRead.ID, model.MessageStatusRead)
assertWebhookMessageStatus(t, db, afterCutoff.ID, model.MessageStatusSent)
var deliveryCount int64
if err := db.Model(&model.DeliveryStatus{}).Where("contact_id = ? AND status = ?", contact.ID, model.MessageStatusRead).Count(&deliveryCount).Error; err != nil {
t.Fatalf("count delivery statuses: %v", err)
}
if deliveryCount != 3 {
t.Fatalf("expected 3 read delivery statuses, got %d", deliveryCount)
}
}
func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) {
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "telegram")
@@ -239,6 +428,37 @@ func listenerSaw(listener *recordingListener, eventType channel.EventType) bool
return false
}
func createWebhookStatusMessage(t *testing.T, db *gorm.DB, inbox model.Inbox, conversation model.Conversation, sourceID string, status model.MessageStatus, createdAt time.Time) model.Message {
t.Helper()
message := model.Message{
ConversationID: conversation.ID,
AccountID: inbox.AccountID,
InboxID: inbox.ID,
Content: sourceID,
ContentType: string(model.MessageContentTypeText),
MessageType: string(model.MessageTypeOutgoing),
Status: string(status),
SourceID: sourceID,
}
message.CreatedAt = createdAt
message.UpdatedAt = createdAt
if err := db.Create(&message).Error; err != nil {
t.Fatalf("create status message: %v", err)
}
return message
}
func assertWebhookMessageStatus(t *testing.T, db *gorm.DB, messageID uint, want model.MessageStatus) {
t.Helper()
var message model.Message
if err := db.First(&message, messageID).Error; err != nil {
t.Fatalf("load message %d: %v", messageID, err)
}
if message.Status != string(want) {
t.Fatalf("message %d expected %s, got %s", messageID, want, message.Status)
}
}
func shopifyHMAC(secret string, body []byte) string {
mac := hmac.New(sha256.New, []byte(secret))
mac.Write(body)
@@ -524,8 +744,8 @@ func TestTwilioDeliveryStatusUpdatesExistingMessage(t *testing.T) {
r.ServeHTTP(w, req)
if w.Code != http.StatusOK {
t.Fatalf("expected 200, got %d", w.Code)
if w.Code != http.StatusNoContent {
t.Fatalf("expected 204, got %d", w.Code)
}
var updated model.Message
if err := db.First(&updated, message.ID).Error; err != nil {
@@ -536,6 +756,76 @@ func TestTwilioDeliveryStatusUpdatesExistingMessage(t *testing.T) {
}
}
func TestTwilioDeliveryStatusQueuesExactRouteWithWorker(t *testing.T) {
gin.SetMode(gin.TestMode)
db := newWebhookLookupTestDB(t)
inbox := seedWebhookInbox(t, db, "twilio_sms")
channelRecord := channelmodel.ChannelTwilioSMS{AccountID: 1, InboxID: inbox.ID, AccountSID: "AC123", PhoneNumber: "+15551234567", MessagingServiceSID: "MG123"}
if err := db.Create(&channelRecord).Error; err != nil {
t.Fatalf("create twilio channel: %v", err)
}
contact := model.Contact{AccountID: inbox.AccountID, Name: "SMS Contact", Identifier: "+15550001111"}
if err := db.Create(&contact).Error; err != nil {
t.Fatalf("create contact: %v", err)
}
conversation := model.Conversation{AccountID: inbox.AccountID, InboxID: inbox.ID, ContactID: contact.ID, Status: string(model.ConversationStatusOpen), ChannelType: inbox.ChannelType, Channel: inbox.ChannelType}
if err := db.Create(&conversation).Error; err != nil {
t.Fatalf("create conversation: %v", err)
}
message := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "out", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusSent), SourceID: "SMFAILED"}
if err := db.Create(&message).Error; err != nil {
t.Fatalf("create message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
h := NewTwilioWebhookHandler(nil, db).WithWorkerPool(wp)
r := gin.New()
r.POST("/twilio/delivery_status", h.HandleTwilioDeliveryStatus)
form := url.Values{}
form.Set("MessagingServiceSid", "MG123")
form.Set("MessageSid", "SMFAILED")
form.Set("MessageStatus", "undelivered")
form.Set("ErrorCode", "30003")
form.Set("ErrorMessage", "Unreachable handset")
req := httptest.NewRequest(http.MethodPost, "/twilio/delivery_status", strings.NewReader(form.Encode()))
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
w := httptest.NewRecorder()
r.ServeHTTP(w, req)
if w.Code != http.StatusNoContent {
t.Fatalf("expected 204, got %d", w.Code)
}
var queued model.BackgroundJob
if err := db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeProviderMessageStatusUpdate, "low", model.BackgroundJobStatusQueued).First(&queued).Error; err != nil {
t.Fatalf("expected queued twilio status job: %v", err)
}
var before model.Message
if err := db.First(&before, message.ID).Error; err != nil {
t.Fatalf("load before message: %v", err)
}
if before.Status != string(model.MessageStatusSent) {
t.Fatalf("expected message to remain sent before worker, got %s", before.Status)
}
processed, err := wp.ProcessOne(t.Context())
if err != nil || !processed {
t.Fatalf("process twilio status job processed=%v err=%v", processed, err)
}
var updated model.Message
if err := db.First(&updated, message.ID).Error; err != nil {
t.Fatalf("load updated message: %v", err)
}
if updated.Status != string(model.MessageStatusFailed) {
t.Fatalf("expected failed status, got %s", updated.Status)
}
attrs := map[string]any{}
if err := json.Unmarshal(updated.ContentAttributes, &attrs); err != nil {
t.Fatalf("unmarshal attrs: %v", err)
}
if attrs["external_error"] != "30003 - Unreachable handset" {
t.Fatalf("expected twilio external error, got %#v", attrs)
}
}
func TestWhatsAppWebhookPersistsIncomingMessage(t *testing.T) {
gin.SetMode(gin.TestMode)
db := newWebhookLookupTestDB(t)
@@ -14,6 +14,7 @@ import (
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/channel/whatsapp"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
// WhatsAppWebhookHandler is a Gin adapter that wraps the WhatsApp
@@ -70,6 +71,13 @@ func NewWhatsAppWebhookHandler(provider *whatsapp.WhatsAppProvider, waWebhook *w
return h
}
func (h *WhatsAppWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *WhatsAppWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetWorkerPool(wp)
}
return h
}
type whatsAppPersisterAdapter struct {
persister *IncomingPersister
}