package webhook import ( "bytes" "context" "crypto/hmac" "crypto/sha256" "encoding/base64" "encoding/hex" "encoding/json" "net/http" "net/http/httptest" "net/url" "strconv" "strings" "testing" "time" "github.com/gin-gonic/gin" "github.com/gochat/gochat/internal/automation" "github.com/gochat/gochat/internal/channel" linechannel "github.com/gochat/gochat/internal/channel/line" channelprovider "github.com/gochat/gochat/internal/channel/provider" tiktokchannel "github.com/gochat/gochat/internal/channel/tiktok" twiliochannel "github.com/gochat/gochat/internal/channel/twilio" 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" ) type recordingListener struct { events []*channel.ChannelEvent } type recordingIncomingSearchIndexer struct { contacts []model.Contact conversations []model.Conversation messages []model.Message } func (r *recordingIncomingSearchIndexer) IndexContact(ctx context.Context, contact *model.Contact) error { r.contacts = append(r.contacts, *contact) return nil } func (r *recordingIncomingSearchIndexer) IndexConversation(ctx context.Context, conversation *model.Conversation) error { r.conversations = append(r.conversations, *conversation) return nil } func (r *recordingIncomingSearchIndexer) IndexMessage(ctx context.Context, message *model.Message) error { r.messages = append(r.messages, *message) return nil } type webhookAutomationDBProvider struct { db *gorm.DB } func (p webhookAutomationDBProvider) DB() *gorm.DB { return p.db } func (l *recordingListener) Name() string { return "recording-listener" } func (l *recordingListener) OnEvent(ctx context.Context, event *channel.ChannelEvent) error { l.events = append(l.events, event) return nil } func TestValidProviderMessageStatusTransition(t *testing.T) { tests := []struct { name string current model.MessageStatus next model.MessageStatus want bool }{ {name: "sent to delivered", current: model.MessageStatusSent, next: model.MessageStatusDelivered, want: true}, {name: "delivered to read", current: model.MessageStatusDelivered, next: model.MessageStatusRead, want: true}, {name: "sent to failed", current: model.MessageStatusSent, next: model.MessageStatusFailed, want: true}, {name: "failed can be recovered by provider", current: model.MessageStatusFailed, next: model.MessageStatusDelivered, want: true}, {name: "duplicate delivered is noop", current: model.MessageStatusDelivered, next: model.MessageStatusDelivered, want: false}, {name: "duplicate failed is noop", current: model.MessageStatusFailed, next: model.MessageStatusFailed, want: false}, {name: "read does not regress to delivered", current: model.MessageStatusRead, next: model.MessageStatusDelivered, want: false}, {name: "delivered does not regress to sent", current: model.MessageStatusDelivered, next: model.MessageStatusSent, want: false}, {name: "invalid next is rejected", current: model.MessageStatusSent, next: model.MessageStatus("queued"), want: false}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { if got := validProviderMessageStatusTransition(tt.current, tt.next); got != tt.want { t.Fatalf("validProviderMessageStatusTransition(%q, %q) = %v, want %v", tt.current, tt.next, got, tt.want) } }) } } func newWebhookLookupTestDB(t *testing.T) *gorm.DB { t.Helper() dsn := "file:" + strings.NewReplacer("/", "_", " ", "_", ":", "_").Replace(t.Name()) + "?mode=memory&cache=shared" db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{}) if err != nil { t.Fatalf("open sqlite: %v", err) } if err := db.AutoMigrate( &model.Inbox{}, &model.Contact{}, &model.ContactInbox{}, &model.Conversation{}, &model.Message{}, &model.DeliveryStatus{}, &model.BackgroundJob{}, &automation.AutomationRule{}, &automation.AutomationExecution{}, &channelmodel.ChannelTelegram{}, &channelmodel.ChannelLINE{}, &channelmodel.ChannelTwilioSMS{}, &channelmodel.ChannelWhatsApp{}, &channelmodel.ChannelTikTok{}, &channelmodel.ChannelFacebook{}, &channelmodel.ChannelInstagram{}, &model.IntegrationHook{}, ); err != nil { t.Fatalf("migrate webhook lookup models: %v", err) } return db } func TestIncomingPersisterTriggersAutomationRuleFromMessageCreated(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") dispatcher := channel.NewDispatcher() automation.RegisterAutomationRuleListener(dispatcher, webhookAutomationDBProvider{db: db}) persister := NewIncomingPersister(db, dispatcher) rule := &automation.AutomationRule{ AccountID: inbox.AccountID, EventName: "message_created", Name: "provider message created", Conditions: automation.Conditions{}, Actions: automation.Actions{}, Active: true, } if err := automation.NewAutomationRuleService(webhookAutomationDBProvider{db: db}).Create(t.Context(), rule); err != nil { t.Fatalf("create automation rule: %v", err) } msg := &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-automation-1", SenderID: "tg-automation-user", SenderName: "Automation User", SenderType: channel.SenderContact, Content: "trigger automation", 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) } var exec automation.AutomationExecution if err := db.Where("rule_id = ? AND conversation_id = ?", rule.ID, result.Conversation.ID).First(&exec).Error; err != nil { t.Fatalf("expected automation execution from provider message_created: %v", err) } if exec.Status != automation.ExecutionStatusSuccess { t.Fatalf("expected success execution, got %s", exec.Status) } } func TestIncomingPersisterUpdatesMessageStatus(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) persister := NewIncomingPersister(db, dispatcher) msg := &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-status-1", SenderID: "tg-user-status", SenderName: "Status User", SenderType: channel.SenderContact, Content: "status me", 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) } if err := persister.UpdateMessageStatus(t.Context(), &inbox, "tg-status-1", model.MessageStatusRead, nil); err != nil { t.Fatalf("update status: %v", err) } var message model.Message if err := db.First(&message, result.Message.ID).Error; err != nil { t.Fatalf("load message: %v", err) } if message.Status != string(model.MessageStatusRead) { t.Fatalf("expected read status, got %s", message.Status) } var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", message.ID, result.Contact.ID).First(&delivery).Error; err != nil { t.Fatalf("expected delivery status: %v", err) } if delivery.Status != model.MessageStatusRead { t.Fatalf("expected delivery read, got %s", delivery.Status) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } } func TestTwilioDeliveryStatusWebhookPersistsFailedStatusAndExternalError(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "twilio_sms") twilioChannel := channelmodel.ChannelTwilioSMS{ AccountID: inbox.AccountID, InboxID: inbox.ID, AccountSID: "ACtwilio", PhoneNumber: "+15551234567", MessagingServiceSID: "MGtwilio", } if err := db.Create(&twilioChannel).Error; err != nil { t.Fatalf("create twilio channel: %v", err) } inbox.ChannelID = twilioChannel.ID if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update inbox channel id: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "Twilio Contact"} 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{ AccountID: inbox.AccountID, InboxID: inbox.ID, ConversationID: conversation.ID, SenderType: "agent", Content: "outbound sms", ContentType: "text", Status: string(model.MessageStatusSent), MessageType: string(model.MessageTypeOutgoing), SourceID: "SMtwilio", } if err := db.Create(&message).Error; err != nil { t.Fatalf("create message: %v", err) } dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) handler := NewTwilioWebhookHandler(nil, db, dispatcher) router := gin.New() router.POST("/webhooks/twilio/status/:phone_number", handler.HandleTwilioDeliveryStatus) form := url.Values{ "MessageSid": {"SMtwilio"}, "MessageStatus": {"undelivered"}, "ErrorCode": {"30007"}, "ErrorMessage": {"Carrier violation"}, } recorder := httptest.NewRecorder() req, _ := http.NewRequest(http.MethodPost, "/webhooks/twilio/status/15551234567", strings.NewReader(form.Encode())) req.Header.Set("Content-Type", "application/x-www-form-urlencoded") router.ServeHTTP(recorder, req) if recorder.Code != http.StatusNoContent { t.Fatalf("expected Twilio no-content ack, got %d body=%s", recorder.Code, recorder.Body.String()) } 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 message status, got %s", updated.Status) } attrs := map[string]any{} if err := json.Unmarshal(updated.ContentAttributes, &attrs); err != nil { t.Fatalf("invalid content attributes: %v", err) } if attrs["external_error"] != "30007 - Carrier violation" { t.Fatalf("expected Twilio external error, got %#v", attrs) } var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", updated.ID, contact.ID).First(&delivery).Error; err != nil { t.Fatalf("expected delivery status: %v", err) } if delivery.Status != model.MessageStatusFailed { t.Fatalf("expected failed delivery status, got %#v", delivery) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } } func TestIncomingPersisterQueuesMessageStatusUpdateWithWorker(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) indexer := &recordingIncomingSearchIndexer{} wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low")) persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer) 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) } persister.SetWorkerPool(wp) 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) } if len(indexer.messages) != 2 || indexer.messages[1].ID != updated.ID || indexer.messages[1].Status != string(model.MessageStatusDelivered) { t.Fatalf("expected delivered message to be indexed, got %#v", indexer.messages) } } 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")) indexer := &recordingIncomingSearchIndexer{} persister := NewIncomingPersister(db).SetSearchIndexer(indexer).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")) indexer := &recordingIncomingSearchIndexer{} persister := NewIncomingPersister(db).SetSearchIndexer(indexer).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) } if len(indexer.messages) != 3 { t.Fatalf("expected 3 updated messages to be indexed, got %#v", indexer.messages) } } func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) indexer := &recordingIncomingSearchIndexer{} persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer) msg := &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-msg-1", SenderID: "tg-user-1", SenderName: "Ada Lovelace", SenderType: channel.SenderContact, Content: "hello", 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) } if result.Contact == nil || result.ContactInbox == nil || result.Conversation == nil || result.Message == nil { t.Fatalf("expected full persistence result: %#v", result) } if result.Message.Content != "hello" || result.Message.SourceID != "tg-msg-1" { t.Fatalf("unexpected message: %#v", result.Message) } if len(indexer.contacts) != 1 || len(indexer.conversations) != 1 || len(indexer.messages) != 1 { t.Fatalf("expected first incoming result indexed, contacts=%#v conversations=%#v messages=%#v", indexer.contacts, indexer.conversations, indexer.messages) } duplicate, err := persister.PersistIncoming(t.Context(), &inbox, msg) if err != nil { t.Fatalf("persist duplicate: %v", err) } if duplicate == nil || !duplicate.Duplicate { t.Fatalf("expected duplicate result, got %#v", duplicate) } if len(indexer.messages) != 1 { t.Fatalf("expected duplicate to skip indexing, got %#v", indexer.messages) } msg.SourceID = "tg-msg-2" msg.Content = "second" second, err := persister.PersistIncoming(t.Context(), &inbox, msg) if err != nil { t.Fatalf("persist second: %v", err) } if second.Contact.ID != result.Contact.ID || second.Conversation.ID != result.Conversation.ID { t.Fatalf("expected contact/conversation reuse: first=%#v second=%#v", result, second) } if len(indexer.messages) != 2 || indexer.messages[1].ID != second.Message.ID { t.Fatalf("expected second incoming message indexed, got %#v", indexer.messages) } var messageCount int64 if err := db.Model(&model.Message{}).Where("inbox_id = ?", inbox.ID).Count(&messageCount).Error; err != nil { t.Fatalf("count messages: %v", err) } if messageCount != 2 { t.Fatalf("expected 2 persisted messages after duplicate skip, got %d", messageCount) } for _, eventType := range []channel.EventType{channel.EventContactCreated, channel.EventConversationCreated, channel.EventConversationOpened, channel.EventMessageCreated, channel.EventMessageIncoming} { if !listenerSaw(listener, eventType) { t.Fatalf("expected event %s, got %#v", eventType, listener.events) } } } func TestIncomingPersisterNormalizesExternalReplyReference(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") persister := NewIncomingPersister(db) original, err := persister.PersistIncoming(t.Context(), &inbox, &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-reply-source-1", SenderID: "tg-reply-user-1", SenderName: "Reply User", SenderType: channel.SenderContact, Content: "original", ContentType: channel.ContentText, InboxID: inbox.ID, AccountID: inbox.AccountID, }) if err != nil { t.Fatalf("persist original message: %v", err) } reply, err := persister.PersistIncoming(t.Context(), &inbox, &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-reply-source-2", SenderID: "tg-reply-user-1", SenderName: "Reply User", SenderType: channel.SenderContact, Content: "reply", ContentType: channel.ContentText, ReplyToID: "tg-reply-source-1", InboxID: inbox.ID, AccountID: inbox.AccountID, }) if err != nil { t.Fatalf("persist reply message: %v", err) } attrs := map[string]interface{}{} if err := json.Unmarshal(reply.Message.ContentAttributes, &attrs); err != nil { t.Fatalf("decode reply content attributes: %v", err) } if got := uint(attrs["in_reply_to"].(float64)); got != original.Message.ID { t.Fatalf("expected internal in_reply_to %d, got %#v", original.Message.ID, attrs["in_reply_to"]) } if got := attrs["in_reply_to_external_id"]; got != original.Message.SourceID { t.Fatalf("expected external reply id %q, got %#v", original.Message.SourceID, got) } unresolved, err := persister.PersistIncoming(t.Context(), &inbox, &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-reply-source-3", SenderID: "tg-reply-user-1", SenderName: "Reply User", SenderType: channel.SenderContact, Content: "unresolved reply", ContentType: channel.ContentText, ReplyToID: "tg-missing-source", InboxID: inbox.ID, AccountID: inbox.AccountID, }) if err != nil { t.Fatalf("persist unresolved reply message: %v", err) } unresolvedAttrs := map[string]interface{}{} if err := json.Unmarshal(unresolved.Message.ContentAttributes, &unresolvedAttrs); err != nil { t.Fatalf("decode unresolved reply content attributes: %v", err) } if got, ok := unresolvedAttrs["in_reply_to"]; !ok || got != nil { t.Fatalf("expected unresolved internal reply id to be null, got %#v", unresolvedAttrs) } if got, ok := unresolvedAttrs["in_reply_to_external_id"]; !ok || got != nil { t.Fatalf("expected unresolved external reply id to be null, got %#v", unresolvedAttrs) } } func TestIncomingPersisterQueuesIncomingMessageWithWorker(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) indexer := &recordingIncomingSearchIndexer{} wp := worker.NewWorkerPool(db) persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer).SetWorkerPool(wp) msg := &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-inbound-job-1", SenderID: "tg-inbound-user-1", SenderName: "Inbound Worker User", SenderType: channel.SenderContact, Content: "persist me later", ContentType: channel.ContentText, InboxID: inbox.ID, AccountID: inbox.AccountID, } result, err := persister.PersistIncoming(t.Context(), &inbox, msg) if err != nil { t.Fatalf("enqueue incoming: %v", err) } if result != nil { t.Fatalf("expected async enqueue result to be nil, got %#v", result) } var queued model.BackgroundJob if err := db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeProviderIncomingMessagePersist, model.DefaultBackgroundJobQueue, model.BackgroundJobStatusQueued).First(&queued).Error; err != nil { t.Fatalf("expected queued incoming job: %v", err) } var beforeCount int64 if err := db.Model(&model.Message{}).Where("source_id = ?", "tg-inbound-job-1").Count(&beforeCount).Error; err != nil { t.Fatalf("count before messages: %v", err) } if beforeCount != 0 { t.Fatalf("expected no message before worker, got %d", beforeCount) } processed, err := wp.ProcessOne(t.Context()) if err != nil || !processed { t.Fatalf("process incoming job processed=%v err=%v", processed, err) } assertPersistedMessage(t, db, inbox.ID, "tg-inbound-job-1", "persist me later") if len(indexer.messages) != 1 || indexer.messages[0].SourceID != "tg-inbound-job-1" { t.Fatalf("expected durable incoming message to be indexed, got %#v", indexer.messages) } for _, eventType := range []channel.EventType{channel.EventContactCreated, channel.EventConversationCreated, channel.EventConversationOpened, channel.EventMessageCreated, channel.EventMessageIncoming} { if !listenerSaw(listener, eventType) { t.Fatalf("expected event %s, got %#v", eventType, listener.events) } } } func TestIncomingPersisterIncomingMessageJobIsIdempotent(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") wp := worker.NewWorkerPool(db) persister := NewIncomingPersister(db).SetWorkerPool(wp) msg := &channel.IncomingMessage{ ChannelType: channel.ChannelTelegram, SourceID: "tg-inbound-idempotent-1", SenderID: "tg-inbound-idempotent-user", SenderName: "Idempotent User", SenderType: channel.SenderContact, Content: "only once", ContentType: channel.ContentText, InboxID: inbox.ID, AccountID: inbox.AccountID, } if _, err := persister.PersistIncoming(t.Context(), &inbox, msg); err != nil { t.Fatalf("enqueue first incoming: %v", err) } if _, err := persister.PersistIncoming(t.Context(), &inbox, msg); err != nil { t.Fatalf("enqueue duplicate incoming: %v", err) } var jobCount int64 if err := db.Model(&model.BackgroundJob{}).Where("job_type = ?", TaskTypeProviderIncomingMessagePersist).Count(&jobCount).Error; err != nil { t.Fatalf("count jobs: %v", err) } if jobCount != 1 { t.Fatalf("expected one idempotent job, got %d", jobCount) } processed, err := wp.ProcessOne(t.Context()) if err != nil || !processed { t.Fatalf("process incoming job processed=%v err=%v", processed, err) } var messageCount int64 if err := db.Model(&model.Message{}).Where("source_id = ?", "tg-inbound-idempotent-1").Count(&messageCount).Error; err != nil { t.Fatalf("count messages: %v", err) } if messageCount != 1 { t.Fatalf("expected one persisted message, got %d", messageCount) } } func listenerSaw(listener *recordingListener, eventType channel.EventType) bool { for _, event := range listener.events { if event.Type == eventType { return true } } 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) return base64.StdEncoding.EncodeToString(mac.Sum(nil)) } func lineSignature(secret string, body []byte) string { mac := hmac.New(sha256.New, []byte(secret)) mac.Write(body) return base64.StdEncoding.EncodeToString(mac.Sum(nil)) } func metaSignature(secret string, body []byte) string { mac := hmac.New(sha256.New, []byte(secret)) mac.Write(body) return "sha256=" + hex.EncodeToString(mac.Sum(nil)) } func tiktokSignature(secret string, timestamp int64, body []byte) string { mac := hmac.New(sha256.New, []byte(secret)) mac.Write([]byte(strconv.FormatInt(timestamp, 10) + "." + string(body))) return "t=" + strconv.FormatInt(timestamp, 10) + ",s=" + hex.EncodeToString(mac.Sum(nil)) } func seedWebhookInbox(t *testing.T, db *gorm.DB, channelType string) model.Inbox { t.Helper() inbox := model.Inbox{ AccountID: 1, Name: channelType + " inbox", ChannelType: channelType, ChannelID: 1, Enabled: true, } if err := db.Create(&inbox).Error; err != nil { t.Fatalf("create inbox: %v", err) } return inbox } func seedInstagramReceiptConversation(t *testing.T, db *gorm.DB, instagramAccountID, contactSourceID string) (model.Inbox, model.Contact, model.Conversation) { t.Helper() inbox := seedWebhookInbox(t, db, "instagram") instagramChannel := channelmodel.ChannelInstagram{ AccountID: inbox.AccountID, InboxID: inbox.ID, InstagramAccountID: instagramAccountID, PageAccessToken: "page-token", ConnectedFBPageID: "fb-page-for-" + instagramAccountID, InstagramAccountName: "support_ig", } if err := db.Create(&instagramChannel).Error; err != nil { t.Fatalf("create instagram channel: %v", err) } inbox.ChannelID = instagramChannel.ID if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update instagram inbox channel id: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "Instagram Contact"} if err := db.Create(&contact).Error; err != nil { t.Fatalf("create contact: %v", err) } contactInbox := model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: contactSourceID, PubsubToken: "pub-" + contactSourceID} 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) } return inbox, contact, conversation } func TestTelegramWebhookLookupInboxByBotToken(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") channel := channelmodel.ChannelTelegram{ AccountID: 1, InboxID: inbox.ID, BotToken: "123:secret-token", BotName: "support_bot", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create telegram channel: %v", err) } h := NewTelegramWebhookHandler(nil, nil, db) found, err := h.lookupInbox("123:secret-token") if err != nil { t.Fatalf("lookup inbox: %v", err) } if found.ID != inbox.ID || found.ChannelType != "telegram" { t.Fatalf("unexpected inbox: id=%d type=%s", found.ID, found.ChannelType) } } func TestTelegramWebhookPersistsIncomingMessage(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "telegram") channelRecord := channelmodel.ChannelTelegram{ AccountID: 1, InboxID: inbox.ID, BotToken: "123:secret-token", BotName: "support_bot", } if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create telegram channel: %v", err) } body := []byte(`{"update_id":1001,"message":{"message_id":2002,"from":{"id":3003,"first_name":"Ada","last_name":"Lovelace","username":"ada"},"chat":{"id":3003,"type":"private"},"date":1710000000,"text":"hello telegram"}}`) h := NewTelegramWebhookHandler(channelprovider.NewTelegramProvider(), nil, db) r := gin.New() r.POST("/webhooks/telegram/:bot_token", h.HandleTelegramWebhook) req := httptest.NewRequest(http.MethodPost, "/webhooks/telegram/123:secret-token", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } var message model.Message if err := db.Where("inbox_id = ? AND source_id = ?", inbox.ID, "2002").First(&message).Error; err != nil { t.Fatalf("expected telegram message persisted: %v", err) } if message.Content != "hello telegram" || message.MessageType != string(model.MessageTypeIncoming) { t.Fatalf("unexpected message: %#v", message) } } func TestFacebookWebhookDeliveryReceiptPersistsDeliveredStatus(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "facebook") facebookChannel := channelmodel.ChannelFacebook{ AccountID: inbox.AccountID, InboxID: inbox.ID, PageID: "page-123", PageAccessToken: "page-token", PageName: "Support Page", } if err := db.Create(&facebookChannel).Error; err != nil { t.Fatalf("create facebook channel: %v", err) } inbox.ChannelID = facebookChannel.ID if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update inbox channel id: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "Facebook Contact"} if err := db.Create(&contact).Error; err != nil { t.Fatalf("create contact: %v", err) } contactInbox := model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: "fb-user-1", PubsubToken: "pub-fb"} 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) } message := createWebhookStatusMessage(t, db, inbox, conversation, "fb-out-1", model.MessageStatusSent, time.Now().Add(-time.Minute).UTC()) dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) handler := NewFacebookWebhookHandler(nil, nil, db, dispatcher) router := gin.New() router.POST("/webhooks/facebook/:page_id", handler.HandleFacebookWebhook) watermark := time.Now().UTC().UnixMilli() body := []byte(`{"object":"page","entry":[{"id":"page-123","time":` + strconv.FormatInt(watermark, 10) + `,"messaging":[{"sender":{"id":"fb-user-1"},"recipient":{"id":"page-123"},"timestamp":` + strconv.FormatInt(watermark, 10) + `,"delivery":{"mids":["fb-out-1"],"watermark":` + strconv.FormatInt(watermark, 10) + `}}]}]}`) recorder := httptest.NewRecorder() req := httptest.NewRequest(http.MethodPost, "/webhooks/facebook/page-123", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") router.ServeHTTP(recorder, req) if recorder.Code != http.StatusOK { t.Fatalf("expected Facebook ack, got %d body=%s", recorder.Code, recorder.Body.String()) } assertWebhookMessageStatus(t, db, message.ID, model.MessageStatusDelivered) var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", message.ID, contact.ID).First(&delivery).Error; err != nil { t.Fatalf("expected delivery status: %v", err) } if delivery.Status != model.MessageStatusDelivered { t.Fatalf("expected delivered delivery status, got %s", delivery.Status) } if delivery.DeliveredAt == nil { t.Fatalf("expected delivered_at timestamp") } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } retry := httptest.NewRecorder() retryReq := httptest.NewRequest(http.MethodPost, "/webhooks/facebook/page-123", bytes.NewReader(body)) retryReq.Header.Set("Content-Type", "application/json") router.ServeHTTP(retry, retryReq) if retry.Code != http.StatusOK { t.Fatalf("expected duplicate Facebook ack, got %d body=%s", retry.Code, retry.Body.String()) } var deliveryCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("message_id = ? AND contact_id = ?", message.ID, contact.ID).Count(&deliveryCount).Error; err != nil { t.Fatalf("count delivery statuses: %v", err) } if deliveryCount != 1 { t.Fatalf("expected duplicate delivery receipt to keep one delivery status row, got %d", deliveryCount) } if len(listener.events) != 1 { t.Fatalf("expected duplicate delivery receipt not to dispatch another status event, got %#v", listener.events) } } func TestFacebookWebhookReadReceiptPersistsReadStatusForContactConversation(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "facebook") facebookChannel := channelmodel.ChannelFacebook{ AccountID: inbox.AccountID, InboxID: inbox.ID, PageID: "page-read-123", PageAccessToken: "page-token", PageName: "Support Page", } if err := db.Create(&facebookChannel).Error; err != nil { t.Fatalf("create facebook channel: %v", err) } inbox.ChannelID = facebookChannel.ID if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update inbox channel id: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "Facebook Read Contact"} if err := db.Create(&contact).Error; err != nil { t.Fatalf("create contact: %v", err) } contactInbox := model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: "fb-user-read", PubsubToken: "pub-fb-read"} 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, "fb-read-sent", model.MessageStatusSent, cutoff.Add(-time.Minute)) beforeDelivered := createWebhookStatusMessage(t, db, inbox, conversation, "fb-read-delivered", model.MessageStatusDelivered, cutoff.Add(-30*time.Second)) dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) handler := NewFacebookWebhookHandler(nil, nil, db, dispatcher) router := gin.New() router.POST("/webhooks/facebook/:page_id", handler.HandleFacebookWebhook) watermark := cutoff.UnixMilli() body := []byte(`{"object":"page","entry":[{"id":"page-read-123","time":` + strconv.FormatInt(watermark, 10) + `,"messaging":[{"sender":{"id":"fb-user-read"},"recipient":{"id":"page-read-123"},"timestamp":` + strconv.FormatInt(watermark, 10) + `,"read":{"watermark":` + strconv.FormatInt(watermark, 10) + `}}]}]}`) recorder := httptest.NewRecorder() req := httptest.NewRequest(http.MethodPost, "/webhooks/facebook/page-read-123", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") router.ServeHTTP(recorder, req) if recorder.Code != http.StatusOK { t.Fatalf("expected Facebook ack, got %d body=%s", recorder.Code, recorder.Body.String()) } assertWebhookMessageStatus(t, db, beforeSent.ID, model.MessageStatusRead) assertWebhookMessageStatus(t, db, beforeDelivered.ID, model.MessageStatusRead) var readCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("contact_id = ? AND status = ?", contact.ID, model.MessageStatusRead).Count(&readCount).Error; err != nil { t.Fatalf("count read delivery statuses: %v", err) } if readCount != 2 { t.Fatalf("expected 2 read delivery statuses, got %d", readCount) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } retry := httptest.NewRecorder() retryReq := httptest.NewRequest(http.MethodPost, "/webhooks/facebook/page-read-123", bytes.NewReader(body)) retryReq.Header.Set("Content-Type", "application/json") router.ServeHTTP(retry, retryReq) if retry.Code != http.StatusOK { t.Fatalf("expected duplicate Facebook read ack, got %d body=%s", retry.Code, retry.Body.String()) } if err := db.Model(&model.DeliveryStatus{}).Where("contact_id = ? AND status = ?", contact.ID, model.MessageStatusRead).Count(&readCount).Error; err != nil { t.Fatalf("count read delivery statuses after retry: %v", err) } if readCount != 2 { t.Fatalf("expected duplicate read receipt to keep two read delivery statuses, got %d", readCount) } if len(listener.events) != 2 { t.Fatalf("expected duplicate read receipt not to dispatch more status events, got %#v", listener.events) } } func TestInstagramWebhookDeliveryReceiptPersistsDeliveredStatus(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("INSTAGRAM_APP_SECRET", "ig-secret") db := newWebhookLookupTestDB(t) inbox, contact, conversation := seedInstagramReceiptConversation(t, db, "ig-account-1", "ig-user-1") message := createWebhookStatusMessage(t, db, inbox, conversation, "ig-out-1", model.MessageStatusSent, time.Now().Add(-time.Minute).UTC()) dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) handler := NewFacebookWebhookHandler(nil, nil, db, dispatcher) router := gin.New() router.POST("/webhooks/instagram", handler.HandleInstagramWebhook) watermark := time.Now().UTC().UnixMilli() body := []byte(`{"object":"instagram","entry":[{"id":"ig-account-1","time":` + strconv.FormatInt(watermark, 10) + `,"messaging":[{"sender":{"id":"ig-user-1"},"recipient":{"id":"ig-account-1"},"timestamp":` + strconv.FormatInt(watermark, 10) + `,"delivery":{"mids":["ig-out-1"],"watermark":` + strconv.FormatInt(watermark, 10) + `}}]}]}`) recorder := httptest.NewRecorder() req := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Hub-Signature-256", metaSignature("ig-secret", body)) router.ServeHTTP(recorder, req) if recorder.Code != http.StatusOK { t.Fatalf("expected Instagram ack, got %d body=%s", recorder.Code, recorder.Body.String()) } assertWebhookMessageStatus(t, db, message.ID, model.MessageStatusDelivered) var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", message.ID, 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 status with timestamp, got %#v", delivery) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } retry := httptest.NewRecorder() retryReq := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) retryReq.Header.Set("Content-Type", "application/json") retryReq.Header.Set("X-Hub-Signature-256", metaSignature("ig-secret", body)) router.ServeHTTP(retry, retryReq) if retry.Code != http.StatusOK { t.Fatalf("expected duplicate Instagram ack, got %d body=%s", retry.Code, retry.Body.String()) } var deliveryCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("message_id = ? AND contact_id = ?", message.ID, contact.ID).Count(&deliveryCount).Error; err != nil { t.Fatalf("count delivery statuses: %v", err) } if deliveryCount != 1 { t.Fatalf("expected duplicate delivery receipt to keep one delivery status row, got %d", deliveryCount) } if len(listener.events) != 1 { t.Fatalf("expected duplicate delivery receipt not to dispatch another status event, got %#v", listener.events) } } func TestInstagramWebhookReadReceiptPersistsReadStatusForContactConversation(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("INSTAGRAM_APP_SECRET", "ig-secret") db := newWebhookLookupTestDB(t) inbox, contact, conversation := seedInstagramReceiptConversation(t, db, "ig-read-account", "ig-read-user") cutoff := time.Now().UTC() beforeSent := createWebhookStatusMessage(t, db, inbox, conversation, "ig-read-sent", model.MessageStatusSent, cutoff.Add(-time.Minute)) beforeDelivered := createWebhookStatusMessage(t, db, inbox, conversation, "ig-read-delivered", model.MessageStatusDelivered, cutoff.Add(-30*time.Second)) dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) handler := NewFacebookWebhookHandler(nil, nil, db, dispatcher) router := gin.New() router.POST("/webhooks/instagram", handler.HandleInstagramWebhook) watermark := cutoff.UnixMilli() body := []byte(`{"object":"instagram","entry":[{"id":"ig-read-account","time":` + strconv.FormatInt(watermark, 10) + `,"messaging":[{"sender":{"id":"ig-read-user"},"recipient":{"id":"ig-read-account"},"timestamp":` + strconv.FormatInt(watermark, 10) + `,"read":{"watermark":` + strconv.FormatInt(watermark, 10) + `}}]}]}`) recorder := httptest.NewRecorder() req := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Hub-Signature-256", metaSignature("ig-secret", body)) router.ServeHTTP(recorder, req) if recorder.Code != http.StatusOK { t.Fatalf("expected Instagram ack, got %d body=%s", recorder.Code, recorder.Body.String()) } assertWebhookMessageStatus(t, db, beforeSent.ID, model.MessageStatusRead) assertWebhookMessageStatus(t, db, beforeDelivered.ID, model.MessageStatusRead) var readCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("contact_id = ? AND status = ?", contact.ID, model.MessageStatusRead).Count(&readCount).Error; err != nil { t.Fatalf("count read delivery statuses: %v", err) } if readCount != 2 { t.Fatalf("expected 2 read delivery statuses, got %d", readCount) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } retry := httptest.NewRecorder() retryReq := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) retryReq.Header.Set("Content-Type", "application/json") retryReq.Header.Set("X-Hub-Signature-256", metaSignature("ig-secret", body)) router.ServeHTTP(retry, retryReq) if retry.Code != http.StatusOK { t.Fatalf("expected duplicate Instagram read ack, got %d body=%s", retry.Code, retry.Body.String()) } if err := db.Model(&model.DeliveryStatus{}).Where("contact_id = ? AND status = ?", contact.ID, model.MessageStatusRead).Count(&readCount).Error; err != nil { t.Fatalf("count read delivery statuses after retry: %v", err) } if readCount != 2 { t.Fatalf("expected duplicate read receipt to keep two read delivery statuses, got %d", readCount) } if len(listener.events) != 2 { t.Fatalf("expected duplicate read receipt not to dispatch more status events, got %#v", listener.events) } } func TestLineWebhookLookupInboxByLineChannelID(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "line") channel := channelmodel.ChannelLINE{ AccountID: 1, InboxID: inbox.ID, ChannelID: "line-channel-1", Name: "LINE OA", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create line channel: %v", err) } h := NewLineWebhookHandler(nil, nil, nil, db) found, err := h.lookupInboxByLineChannelID("line-channel-1") if err != nil { t.Fatalf("lookup inbox: %v", err) } if found.ID != inbox.ID || found.ChannelType != "line" { t.Fatalf("unexpected inbox: id=%d type=%s", found.ID, found.ChannelType) } } func TestLineWebhookPersistsIncomingMessage(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "line") inbox.ChannelConfig = `{"channel_secret":"line-secret"}` if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update line inbox config: %v", err) } channelRecord := channelmodel.ChannelLINE{AccountID: 1, InboxID: inbox.ID, ChannelID: "line-channel-1", Name: "LINE OA"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create line channel: %v", err) } lineRepo := linechannel.NewRepository(db) lineService := linechannel.NewLineService(lineRepo) linePipeline := linechannel.NewIncomingProcessor(lineService) h := NewLineWebhookHandler(nil, linePipeline, lineService, db) r := gin.New() r.POST("/webhooks/line/:line_channel_id", h.HandleLineWebhook) body := []byte(`{"destination":"line-channel-1","events":[{"type":"message","replyToken":"reply-1","timestamp":1710000000000,"source":{"type":"user","userId":"line-user-1"},"message":{"type":"text","id":"line-msg-1","text":"hello line"}}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/line/line-channel-1", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Line-Signature", lineSignature("line-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "line-msg-1", "hello line") } func TestLineWebhookRejectsMissingSignatureWhenSecretConfigured(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "line") inbox.ChannelConfig = `{"channel_secret":"line-secret"}` if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update line inbox config: %v", err) } channelRecord := channelmodel.ChannelLINE{AccountID: 1, InboxID: inbox.ID, ChannelID: "line-channel-1", Name: "LINE OA"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create line channel: %v", err) } lineRepo := linechannel.NewRepository(db) lineService := linechannel.NewLineService(lineRepo) linePipeline := linechannel.NewIncomingProcessor(lineService) h := NewLineWebhookHandler(nil, linePipeline, lineService, db) r := gin.New() r.POST("/webhooks/line/:line_channel_id", h.HandleLineWebhook) body := []byte(`{"destination":"line-channel-1","events":[{"type":"message","replyToken":"reply-1","timestamp":1710000000000,"source":{"type":"user","userId":"line-user-1"},"message":{"type":"text","id":"line-msg-missing-sig","text":"hello line"}}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/line/line-channel-1", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d body=%s", w.Code, w.Body.String()) } var count int64 if err := db.Model(&model.Message{}).Where("inbox_id = ? AND source_id = ?", inbox.ID, "line-msg-missing-sig").Count(&count).Error; err != nil { t.Fatalf("count message: %v", err) } if count != 0 { t.Fatalf("expected no persisted message, got %d", count) } } func TestLineWebhookRejectsInvalidSignatureWhenSecretConfigured(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "line") inbox.ChannelConfig = `{"channel_secret":"line-secret"}` if err := db.Save(&inbox).Error; err != nil { t.Fatalf("update line inbox config: %v", err) } channelRecord := channelmodel.ChannelLINE{AccountID: 1, InboxID: inbox.ID, ChannelID: "line-channel-1", Name: "LINE OA"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create line channel: %v", err) } lineRepo := linechannel.NewRepository(db) lineService := linechannel.NewLineService(lineRepo) linePipeline := linechannel.NewIncomingProcessor(lineService) h := NewLineWebhookHandler(nil, linePipeline, lineService, db) r := gin.New() r.POST("/webhooks/line/:line_channel_id", h.HandleLineWebhook) body := []byte(`{"destination":"line-channel-1","events":[{"type":"message","replyToken":"reply-1","timestamp":1710000000000,"source":{"type":"user","userId":"line-user-1"},"message":{"type":"text","id":"line-msg-bad-sig","text":"hello line"}}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/line/line-channel-1", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Line-Signature", lineSignature("wrong-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d body=%s", w.Code, w.Body.String()) } var count int64 if err := db.Model(&model.Message{}).Where("inbox_id = ? AND source_id = ?", inbox.ID, "line-msg-bad-sig").Count(&count).Error; err != nil { t.Fatalf("count message: %v", err) } if count != 0 { t.Fatalf("expected no persisted message, got %d", count) } } func TestTwilioWebhookLookupInboxByPhoneNumber(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "twilio_sms") channel := channelmodel.ChannelTwilioSMS{ AccountID: 1, InboxID: inbox.ID, AccountSID: "AC123", PhoneNumber: "+15551234567", MessagingServiceSID: "MG123", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create twilio channel: %v", err) } h := NewTwilioWebhookHandler(nil, db) found, err := h.lookupInboxByPhoneNumber("+15551234567") if err != nil { t.Fatalf("lookup inbox: %v", err) } if found.ID != inbox.ID || found.ChannelType != "twilio_sms" { t.Fatalf("unexpected inbox: id=%d type=%s", found.ID, found.ChannelType) } } func TestTwilioWebhookPersistsIncomingMessage(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"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create twilio channel: %v", err) } twilioRepo := twiliochannel.NewRepository(db) twilioService := twiliochannel.NewTwilioService(twilioRepo) twilioPipeline := twiliochannel.NewIncomingProcessor(twilioService) twilioWebhook := twiliochannel.NewWebhookHandler(twilioPipeline, twilioService) h := NewTwilioWebhookHandler(twilioWebhook, db) r := gin.New() r.POST("/webhooks/sms/:phone_number", h.HandleTwilioInboundSMS) form := url.Values{} form.Set("MessageSid", "SMIN1") form.Set("AccountSid", "AC123") form.Set("From", "+15550002222") form.Set("To", "+15551234567") form.Set("Body", "hello sms") req := httptest.NewRequest(http.MethodPost, "/webhooks/sms/+15551234567", 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.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "SMIN1", "hello sms") } func TestTwilioCallbackExactRoutePersistsIncomingMessage(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) } twilioRepo := twiliochannel.NewRepository(db) twilioService := twiliochannel.NewTwilioService(twilioRepo) twilioPipeline := twiliochannel.NewIncomingProcessor(twilioService) twilioWebhook := twiliochannel.NewWebhookHandler(twilioPipeline, twilioService) h := NewTwilioWebhookHandler(twilioWebhook, db) r := gin.New() r.POST("/twilio/callback", h.HandleTwilioCallback) form := url.Values{} form.Set("MessageSid", "SMROOT1") form.Set("AccountSid", "AC123") form.Set("From", "+15550002222") form.Set("To", "15551234567") form.Set("Body", "hello root callback") req := httptest.NewRequest(http.MethodPost, "/twilio/callback", 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 body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "SMROOT1", "hello root callback") } func TestTwilioCallbackExactRouteFallsBackToMessagingServiceSid(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) } twilioRepo := twiliochannel.NewRepository(db) twilioService := twiliochannel.NewTwilioService(twilioRepo) twilioPipeline := twiliochannel.NewIncomingProcessor(twilioService) twilioWebhook := twiliochannel.NewWebhookHandler(twilioPipeline, twilioService) h := NewTwilioWebhookHandler(twilioWebhook, db) r := gin.New() r.POST("/twilio/callback", h.HandleTwilioCallback) form := url.Values{} form.Set("MessageSid", "SMROOTMG1") form.Set("MessagingServiceSid", "MG123") form.Set("From", "+15550002222") form.Set("Body", "hello service callback") req := httptest.NewRequest(http.MethodPost, "/twilio/callback", 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 body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "SMROOTMG1", "hello service callback") } func TestTwilioInboundSMSQueuesIncomingMessageWithWorker(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"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create twilio channel: %v", err) } twilioRepo := twiliochannel.NewRepository(db) twilioService := twiliochannel.NewTwilioService(twilioRepo) twilioPipeline := twiliochannel.NewIncomingProcessor(twilioService) twilioWebhook := twiliochannel.NewWebhookHandler(twilioPipeline, twilioService) wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low")) h := NewTwilioWebhookHandler(twilioWebhook, db).WithWorkerPool(wp) r := gin.New() r.POST("/webhooks/sms/:phone_number", h.HandleTwilioInboundSMS) form := url.Values{} form.Set("MessageSid", "SMINASYNC1") form.Set("AccountSid", "AC123") form.Set("From", "+15550002222") form.Set("To", "+15551234567") form.Set("Body", "queued sms") req := httptest.NewRequest(http.MethodPost, "/webhooks/sms/+15551234567", 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.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } var queued model.BackgroundJob if err := db.Where("job_type = ? AND queue = ? AND status = ?", TaskTypeProviderIncomingMessagePersist, "low", model.BackgroundJobStatusQueued).First(&queued).Error; err != nil { t.Fatalf("expected queued twilio incoming job: %v", err) } var beforeCount int64 if err := db.Model(&model.Message{}).Where("source_id = ?", "SMINASYNC1").Count(&beforeCount).Error; err != nil { t.Fatalf("count before messages: %v", err) } if beforeCount != 0 { t.Fatalf("expected no message before worker, got %d", beforeCount) } processed, err := wp.ProcessOne(t.Context()) if err != nil || !processed { t.Fatalf("process twilio incoming job processed=%v err=%v", processed, err) } assertPersistedMessage(t, db, inbox.ID, "SMINASYNC1", "queued sms") } func TestTwilioDeliveryStatusPersistsDeliveredAndReadStatuses(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"} 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) } contactInbox := model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, SourceID: "+15550001111", PubsubToken: "pub-twilio"} 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) } deliveredMessage := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "delivered", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusSent), SourceID: "SMDELIVERED"} readMessage := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "read", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusDelivered), SourceID: "SMREAD"} if err := db.Create(&deliveredMessage).Error; err != nil { t.Fatalf("create delivered message: %v", err) } if err := db.Create(&readMessage).Error; err != nil { t.Fatalf("create read message: %v", err) } dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) h := NewTwilioWebhookHandler(nil, db, dispatcher) r := gin.New() r.POST("/webhooks/twilio/status/:phone_number", h.HandleTwilioDeliveryStatus) for _, form := range []string{ "MessageSid=SMDELIVERED&MessageStatus=delivered", "MessageSid=SMREAD&MessageStatus=read", } { req := httptest.NewRequest(http.MethodPost, "/webhooks/twilio/status/+15551234567", strings.NewReader(form)) 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) } } assertWebhookMessageStatus(t, db, deliveredMessage.ID, model.MessageStatusDelivered) assertWebhookMessageStatus(t, db, readMessage.ID, model.MessageStatusRead) var delivered model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", deliveredMessage.ID, contact.ID).First(&delivered).Error; err != nil { t.Fatalf("expected delivered delivery status: %v", err) } if delivered.Status != model.MessageStatusDelivered { t.Fatalf("expected delivered delivery status, got %#v", delivered) } var read model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", readMessage.ID, contact.ID).First(&read).Error; err != nil { t.Fatalf("expected read delivery status: %v", err) } if read.Status != model.MessageStatusRead { t.Fatalf("expected read delivery status, got %#v", read) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } duplicate := httptest.NewRecorder() duplicateReq := httptest.NewRequest(http.MethodPost, "/webhooks/twilio/status/+15551234567", strings.NewReader("MessageSid=SMDELIVERED&MessageStatus=delivered")) duplicateReq.Header.Set("Content-Type", "application/x-www-form-urlencoded") r.ServeHTTP(duplicate, duplicateReq) if duplicate.Code != http.StatusNoContent { t.Fatalf("expected duplicate callback 204, got %d", duplicate.Code) } var deliveredCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("message_id = ? AND contact_id = ?", deliveredMessage.ID, contact.ID).Count(&deliveredCount).Error; err != nil { t.Fatalf("count delivered status rows: %v", err) } if deliveredCount != 1 { t.Fatalf("expected duplicate callback to keep one delivery status row, got %d", deliveredCount) } if len(listener.events) != 2 { t.Fatalf("expected duplicate callback not to dispatch another status event, got %#v", listener.events) } } 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) inbox := seedWebhookInbox(t, db, "whatsapp") waChannel := channelmodel.ChannelWhatsApp{AccountID: 1, InboxID: inbox.ID, PhoneNumber: "+15551230000", PhoneNumberID: "phone-id-1", AccessToken: "token", Provider: "whatsapp_cloud", ProviderConfig: `{"app_secret":"wa-secret"}`, WebhookVerifyToken: "verify-token"} if err := db.Create(&waChannel).Error; err != nil { t.Fatalf("create whatsapp channel: %v", err) } waRepo := whatsappchannel.NewRepository(db) waService := whatsappchannel.NewWhatsAppService(waRepo) waPipeline := whatsappchannel.NewIncomingPipeline(waService) waProvider := whatsappchannel.NewWhatsAppProvider(waService, waRepo, waPipeline) waWebhook := whatsappchannel.NewWebhookHandler(waProvider) h := NewWhatsAppWebhookHandler(waProvider, waWebhook, db) r := gin.New() r.POST("/webhooks/whatsapp/:phone_number", h.HandleWhatsAppWebhook) body := []byte(`{"object":"whatsapp_business_account","entry":[{"id":"waba-1","changes":[{"field":"messages","value":{"messaging_product":"whatsapp","metadata":{"display_phone_number":"+15551230000","phone_number_id":"phone-id-1"},"contacts":[{"wa_id":"15550001111","profile":{"name":"WhatsApp User"}}],"messages":[{"from":"15550001111","id":"wamid-1","timestamp":"1710000000","type":"text","text":{"body":"hello whatsapp"}}]}}]}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/whatsapp/+15551230000", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Hub-Signature-256", metaSignature("wa-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "wamid-1", "hello whatsapp") } func TestWhatsAppWebhookPersistsDeliveryStatusAndFailureError(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "whatsapp") waChannel := channelmodel.ChannelWhatsApp{AccountID: inbox.AccountID, InboxID: inbox.ID, PhoneNumber: "+15551230000", PhoneNumberID: "phone-id-1", AccessToken: "token", Provider: "whatsapp_cloud", ProviderConfig: `{"app_secret":"wa-secret"}`, WebhookVerifyToken: "verify-token"} if err := db.Create(&waChannel).Error; err != nil { t.Fatalf("create whatsapp channel: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "WhatsApp 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: "wamid-out-1"} if err := db.Create(&message).Error; err != nil { t.Fatalf("create message: %v", err) } waRepo := whatsappchannel.NewRepository(db) waService := whatsappchannel.NewWhatsAppService(waRepo) waPipeline := whatsappchannel.NewIncomingPipeline(waService) waProvider := whatsappchannel.NewWhatsAppProvider(waService, waRepo, waPipeline) waWebhook := whatsappchannel.NewWebhookHandler(waProvider) h := NewWhatsAppWebhookHandler(waProvider, waWebhook, db) r := gin.New() r.POST("/webhooks/whatsapp/:phone_number", h.HandleWhatsAppWebhook) body := []byte(`{"object":"whatsapp_business_account","entry":[{"id":"waba-1","changes":[{"field":"messages","value":{"messaging_product":"whatsapp","metadata":{"display_phone_number":"+15551230000","phone_number_id":"phone-id-1"},"statuses":[{"id":"wamid-out-1","status":"failed","timestamp":"1710000100","recipient_id":"15550001111","errors":[{"code":131047,"title":"Re-engagement message","message":"Message failed"}]}]}}]}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/whatsapp/+15551230000", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Hub-Signature-256", metaSignature("wa-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } 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"] != "131047 - Re-engagement message" { t.Fatalf("expected whatsapp external error, got %#v", attrs) } var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", updated.ID, contact.ID).First(&delivery).Error; err != nil { t.Fatalf("expected delivery status: %v", err) } if delivery.Status != model.MessageStatusFailed { t.Fatalf("expected delivery failed, got %s", delivery.Status) } } func TestWhatsAppWebhookPersistsDeliveredAndReadStatuses(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "whatsapp") waChannel := channelmodel.ChannelWhatsApp{AccountID: inbox.AccountID, InboxID: inbox.ID, PhoneNumber: "+15551230000", PhoneNumberID: "phone-id-1", AccessToken: "token", Provider: "whatsapp_cloud", ProviderConfig: `{"app_secret":"wa-secret"}`, WebhookVerifyToken: "verify-token"} if err := db.Create(&waChannel).Error; err != nil { t.Fatalf("create whatsapp channel: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "WhatsApp Status Contact", Identifier: "15550002222"} 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) } deliveredMessage := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "delivered", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusSent), SourceID: "wamid-delivered-1"} readMessage := model.Message{ConversationID: conversation.ID, AccountID: inbox.AccountID, InboxID: inbox.ID, SenderID: &contact.ID, SenderType: string(model.SenderTypeContact), Content: "read", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusDelivered), SourceID: "wamid-read-1"} if err := db.Create(&deliveredMessage).Error; err != nil { t.Fatalf("create delivered message: %v", err) } if err := db.Create(&readMessage).Error; err != nil { t.Fatalf("create read message: %v", err) } waRepo := whatsappchannel.NewRepository(db) waService := whatsappchannel.NewWhatsAppService(waRepo) waPipeline := whatsappchannel.NewIncomingPipeline(waService) waProvider := whatsappchannel.NewWhatsAppProvider(waService, waRepo, waPipeline) waWebhook := whatsappchannel.NewWebhookHandler(waProvider) dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) h := NewWhatsAppWebhookHandler(waProvider, waWebhook, db, dispatcher) r := gin.New() r.POST("/webhooks/whatsapp/:phone_number", h.HandleWhatsAppWebhook) body := []byte(`{"object":"whatsapp_business_account","entry":[{"id":"waba-1","changes":[{"field":"messages","value":{"messaging_product":"whatsapp","metadata":{"display_phone_number":"+15551230000","phone_number_id":"phone-id-1"},"statuses":[{"id":"wamid-delivered-1","status":"delivered","timestamp":"1710000200","recipient_id":"15550002222"},{"id":"wamid-read-1","status":"read","timestamp":"1710000300","recipient_id":"15550002222"}]}}]}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/whatsapp/+15551230000", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Hub-Signature-256", metaSignature("wa-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertWebhookMessageStatus(t, db, deliveredMessage.ID, model.MessageStatusDelivered) assertWebhookMessageStatus(t, db, readMessage.ID, model.MessageStatusRead) var delivered model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", deliveredMessage.ID, contact.ID).First(&delivered).Error; err != nil { t.Fatalf("expected delivered delivery status: %v", err) } if delivered.Status != model.MessageStatusDelivered || delivered.DeliveredAt == nil { t.Fatalf("expected delivered status timestamp, got %#v", delivered) } var read model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", readMessage.ID, contact.ID).First(&read).Error; err != nil { t.Fatalf("expected read delivery status: %v", err) } if read.Status != model.MessageStatusRead || read.ReadAt == nil { t.Fatalf("expected read status timestamp, got %#v", read) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } duplicateBody := []byte(`{"object":"whatsapp_business_account","entry":[{"id":"waba-1","changes":[{"field":"messages","value":{"messaging_product":"whatsapp","metadata":{"display_phone_number":"+15551230000","phone_number_id":"phone-id-1"},"statuses":[{"id":"wamid-delivered-1","status":"delivered","timestamp":"1710000400","recipient_id":"15550002222"}]}}]}]}`) duplicateReq := httptest.NewRequest(http.MethodPost, "/webhooks/whatsapp/+15551230000", bytes.NewReader(duplicateBody)) duplicateReq.Header.Set("Content-Type", "application/json") duplicateReq.Header.Set("X-Hub-Signature-256", metaSignature("wa-secret", duplicateBody)) duplicate := httptest.NewRecorder() r.ServeHTTP(duplicate, duplicateReq) if duplicate.Code != http.StatusOK { t.Fatalf("expected duplicate callback 200, got %d body=%s", duplicate.Code, duplicate.Body.String()) } var deliveredCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("message_id = ? AND contact_id = ?", deliveredMessage.ID, contact.ID).Count(&deliveredCount).Error; err != nil { t.Fatalf("count delivered status rows: %v", err) } if deliveredCount != 1 { t.Fatalf("expected duplicate callback to keep one delivery status row, got %d", deliveredCount) } if len(listener.events) != 2 { t.Fatalf("expected duplicate callback not to dispatch another status event, got %#v", listener.events) } } func TestWhatsAppWebhookVerificationEchoesChallenge(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "whatsapp") waChannel := channelmodel.ChannelWhatsApp{AccountID: 1, InboxID: inbox.ID, PhoneNumber: "+15551230000", PhoneNumberID: "phone-id-1", AccessToken: "token", Provider: "whatsapp_cloud", WebhookVerifyToken: "verify-token"} if err := db.Create(&waChannel).Error; err != nil { t.Fatalf("create whatsapp channel: %v", err) } waRepo := whatsappchannel.NewRepository(db) waService := whatsappchannel.NewWhatsAppService(waRepo) waPipeline := whatsappchannel.NewIncomingPipeline(waService) waProvider := whatsappchannel.NewWhatsAppProvider(waService, waRepo, waPipeline) waWebhook := whatsappchannel.NewWebhookHandler(waProvider) h := NewWhatsAppWebhookHandler(waProvider, waWebhook, db) r := gin.New() r.GET("/webhooks/whatsapp/:phone_number", h.HandleWhatsAppVerification) req := httptest.NewRequest(http.MethodGet, "/webhooks/whatsapp/+15551230000?hub.mode=subscribe&hub.verify_token=verify-token&hub.challenge=challenge-wa", nil) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } if w.Body.String() != "challenge-wa" { t.Fatalf("unexpected challenge body: %q", w.Body.String()) } } func TestWhatsAppCloudWebhookRejectsMissingSignature(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "whatsapp") waChannel := channelmodel.ChannelWhatsApp{AccountID: 1, InboxID: inbox.ID, PhoneNumber: "+15551230000", PhoneNumberID: "phone-id-1", AccessToken: "token", Provider: "whatsapp_cloud", ProviderConfig: `{"app_secret":"wa-secret"}`, WebhookVerifyToken: "verify-token"} if err := db.Create(&waChannel).Error; err != nil { t.Fatalf("create whatsapp channel: %v", err) } waRepo := whatsappchannel.NewRepository(db) waService := whatsappchannel.NewWhatsAppService(waRepo) waPipeline := whatsappchannel.NewIncomingPipeline(waService) waProvider := whatsappchannel.NewWhatsAppProvider(waService, waRepo, waPipeline) waWebhook := whatsappchannel.NewWebhookHandler(waProvider) h := NewWhatsAppWebhookHandler(waProvider, waWebhook, db) r := gin.New() r.POST("/webhooks/whatsapp/:phone_number", h.HandleWhatsAppWebhook) body := []byte(`{"object":"whatsapp_business_account","entry":[{"id":"waba-1","changes":[{"field":"messages","value":{"messaging_product":"whatsapp","metadata":{"display_phone_number":"+15551230000","phone_number_id":"phone-id-1"},"contacts":[{"wa_id":"15550001111","profile":{"name":"WhatsApp User"}}],"messages":[{"from":"15550001111","id":"wamid-missing-sig","timestamp":"1710000000","type":"text","text":{"body":"hello whatsapp"}}]}}]}]}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/whatsapp/+15551230000", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d body=%s", w.Code, w.Body.String()) } var count int64 if err := db.Model(&model.Message{}).Where("inbox_id = ? AND source_id = ?", inbox.ID, "wamid-missing-sig").Count(&count).Error; err != nil { t.Fatalf("count message: %v", err) } if count != 0 { t.Fatalf("expected no persisted message, got %d", count) } } func TestTikTokWebhookLookupInboxByBusinessIDAndPayloadExtractor(t *testing.T) { db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "tiktok") channel := channelmodel.ChannelTikTok{ AccountID: 1, InboxID: inbox.ID, TikTokBusinessID: "biz-123", WebhookVerifyToken: "verify-token", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create tiktok channel: %v", err) } h := NewTikTokWebhookHandler(nil, nil, db) found, err := h.lookupInboxByBusinessID("biz-123") if err != nil { t.Fatalf("lookup inbox: %v", err) } if found.ID != inbox.ID || found.ChannelType != "tiktok" { t.Fatalf("unexpected inbox: id=%d type=%s", found.ID, found.ChannelType) } if got := extractTikTokBusinessID([]byte(`{"data":{"business_id":"biz-123"}}`)); got != "biz-123" { t.Fatalf("unexpected extracted business id: %s", got) } } func TestTikTokWebhookPersistsIncomingMessage(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("TIKTOK_APP_SECRET", "tiktok-secret") db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "tiktok") channelRecord := channelmodel.ChannelTikTok{AccountID: 1, InboxID: inbox.ID, TikTokBusinessID: "biz-123", WebhookVerifyToken: "verify-token"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create tiktok channel: %v", err) } ttRepo := tiktokchannel.NewRepository(db) ttService := tiktokchannel.NewTikTokService(ttRepo) ttPipeline := tiktokchannel.NewIncomingProcessor(ttService, ttRepo) ttWebhook := tiktokchannel.NewWebhookHandler(ttService, ttPipeline) h := NewTikTokWebhookHandler(ttWebhook, ttPipeline, db) r := gin.New() r.POST("/webhooks/tiktok", h.HandleTikTokWebhook) body := []byte(`{"type":"message.received","timestamp":1710000000,"biz_id":"biz-123","data":{"message_id":"tt-msg-1","from_user_id":"tt-user-1","to_user_id":"biz-123","content_type":"text","content":"hello tiktok","timestamp":1710000000,"conversation_id":"tt-conv-1"}}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/tiktok", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Tiktok-Signature", tiktokSignature("tiktok-secret", time.Now().Unix(), body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "tt-msg-1", "hello tiktok") } func TestTikTokWebhookReadReceiptUpdatesMessageStatus(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("TIKTOK_APP_SECRET", "tiktok-secret") db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "tiktok") channelRecord := channelmodel.ChannelTikTok{AccountID: 1, InboxID: inbox.ID, TikTokBusinessID: "biz-123", WebhookVerifyToken: "verify-token"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create tiktok channel: %v", err) } contact := model.Contact{AccountID: inbox.AccountID, Name: "TikTok Contact"} 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: "outbound tiktok", ContentType: string(model.MessageContentTypeText), MessageType: string(model.MessageTypeOutgoing), Status: string(model.MessageStatusDelivered), SourceID: "tt-out-1"} if err := db.Create(&message).Error; err != nil { t.Fatalf("create message: %v", err) } dispatcher := channel.NewDispatcher() listener := &recordingListener{} dispatcher.Register(listener) ttRepo := tiktokchannel.NewRepository(db) ttService := tiktokchannel.NewTikTokService(ttRepo) ttPipeline := tiktokchannel.NewIncomingProcessor(ttService, ttRepo) ttWebhook := tiktokchannel.NewWebhookHandler(ttService, ttPipeline) h := NewTikTokWebhookHandler(ttWebhook, ttPipeline, db, dispatcher) r := gin.New() r.POST("/webhooks/tiktok", h.HandleTikTokWebhook) body := []byte(`{"type":"message.read","timestamp":1710000001,"biz_id":"biz-123","data":{"message_id":"tt-out-1","from_user_id":"tt-user-1","timestamp":1710000001}}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/tiktok", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Tiktok-Signature", tiktokSignature("tiktok-secret", time.Now().Unix(), body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } 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.MessageStatusRead) { t.Fatalf("expected read message status, got %s", updated.Status) } var delivery model.DeliveryStatus if err := db.Where("message_id = ? AND contact_id = ?", updated.ID, contact.ID).First(&delivery).Error; err != nil { t.Fatalf("expected delivery status: %v", err) } if delivery.Status != model.MessageStatusRead { t.Fatalf("expected read delivery status, got %#v", delivery) } if !listenerSaw(listener, channel.EventMessageStatusUpdated) { t.Fatalf("expected message.status_updated event, got %#v", listener.events) } retry := httptest.NewRecorder() retryReq := httptest.NewRequest(http.MethodPost, "/webhooks/tiktok", bytes.NewReader(body)) retryReq.Header.Set("Content-Type", "application/json") retryReq.Header.Set("Tiktok-Signature", tiktokSignature("tiktok-secret", time.Now().Unix(), body)) r.ServeHTTP(retry, retryReq) if retry.Code != http.StatusOK { t.Fatalf("expected duplicate TikTok read ack, got %d body=%s", retry.Code, retry.Body.String()) } var deliveryCount int64 if err := db.Model(&model.DeliveryStatus{}).Where("message_id = ? AND contact_id = ?", updated.ID, contact.ID).Count(&deliveryCount).Error; err != nil { t.Fatalf("count delivery statuses after retry: %v", err) } if deliveryCount != 1 { t.Fatalf("expected duplicate read receipt to keep one delivery status row, got %d", deliveryCount) } if len(listener.events) != 1 { t.Fatalf("expected duplicate read receipt not to dispatch another status event, got %#v", listener.events) } } func TestTikTokWebhookRejectsInvalidSignature(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("TIKTOK_APP_SECRET", "tiktok-secret") db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "tiktok") channelRecord := channelmodel.ChannelTikTok{AccountID: 1, InboxID: inbox.ID, TikTokBusinessID: "biz-123", WebhookVerifyToken: "verify-token"} if err := db.Create(&channelRecord).Error; err != nil { t.Fatalf("create tiktok channel: %v", err) } ttRepo := tiktokchannel.NewRepository(db) ttService := tiktokchannel.NewTikTokService(ttRepo) ttPipeline := tiktokchannel.NewIncomingProcessor(ttService, ttRepo) ttWebhook := tiktokchannel.NewWebhookHandler(ttService, ttPipeline) h := NewTikTokWebhookHandler(ttWebhook, ttPipeline, db) r := gin.New() r.POST("/webhooks/tiktok", h.HandleTikTokWebhook) body := []byte(`{"type":"message.received","timestamp":1710000000,"biz_id":"biz-123","data":{"message_id":"tt-msg-invalid","from_user_id":"tt-user-1","to_user_id":"biz-123","content_type":"text","content":"hello tiktok","timestamp":1710000000,"conversation_id":"tt-conv-1"}}`) req := httptest.NewRequest(http.MethodPost, "/webhooks/tiktok", bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Tiktok-Signature", "t="+strconv.FormatInt(time.Now().Unix(), 10)+",s=bad") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d body=%s", w.Code, w.Body.String()) } var count int64 if err := db.Model(&model.Message{}).Where("inbox_id = ? AND source_id = ?", inbox.ID, "tt-msg-invalid").Count(&count).Error; err != nil { t.Fatalf("count message: %v", err) } if count != 0 { t.Fatalf("expected no persisted message, got %d", count) } } func TestTikTokWebhookRejectsStaleSignature(t *testing.T) { t.Setenv("TIKTOK_APP_SECRET", "tiktok-secret") body := []byte(`{"type":"message.received","biz_id":"biz-123"}`) timestamp := time.Now().Add(-10 * time.Second).Unix() if err := verifyTikTokSignature(tiktokSignature("tiktok-secret", timestamp, body), body, time.Now()); err == nil { t.Fatal("expected stale signature rejection") } } func TestShopifyWebhookShopRedactDeletesMatchingHook(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) settings, _ := json.Marshal(model.ShopifySettings{ShopDomain: "store.myshopify.com"}) hook := model.IntegrationHook{ AccountID: 1, HookType: model.HookTypeShopify, Status: model.HookStatusActive, AccessToken: "access-token", Settings: settings, } if err := db.Create(&hook).Error; err != nil { t.Fatalf("create shopify hook: %v", err) } body := []byte(`{"shop_domain":"store.myshopify.com"}`) h := NewShopifyWebhookHandler(db, "client-secret") r := gin.New() r.POST("/webhooks/shopify", h.HandleShopifyWebhook) req := httptest.NewRequest(http.MethodPost, "/webhooks/shopify", bytes.NewReader(body)) req.Header.Set("X-Shopify-Hmac-SHA256", shopifyHMAC("client-secret", body)) req.Header.Set("X-Shopify-Topic", "shop/redact") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d", w.Code) } var count int64 if err := db.Model(&model.IntegrationHook{}).Where("id = ?", hook.ID).Count(&count).Error; err != nil { t.Fatalf("count hook: %v", err) } if count != 0 { t.Fatalf("expected shopify hook to be deleted, count=%d", count) } } func TestShopifyWebhookRejectsInvalidHMAC(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) h := NewShopifyWebhookHandler(db, "client-secret") r := gin.New() r.POST("/webhooks/shopify", h.HandleShopifyWebhook) req := httptest.NewRequest(http.MethodPost, "/webhooks/shopify", bytes.NewReader([]byte(`{}`))) req.Header.Set("X-Shopify-Hmac-SHA256", "invalid") w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d", w.Code) } } func TestInstagramWebhookVerificationUsesChatwootGlobalTokens(t *testing.T) { gin.SetMode(gin.TestMode) t.Setenv("INSTAGRAM_VERIFY_TOKEN", "verify-me") h := NewFacebookWebhookHandler(nil, nil, nil) r := gin.New() r.GET("/webhooks/instagram", h.HandleInstagramVerification) w := httptest.NewRecorder() req := httptest.NewRequest(http.MethodGet, "/webhooks/instagram?hub.verify_token=verify-me&hub.challenge=challenge-1", nil) r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d", w.Code) } if w.Body.String() != "challenge-1" { t.Fatalf("unexpected challenge body: %q", w.Body.String()) } } func TestInstagramWebhookEventsVerifySignatureAndResolveInbox(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "instagram") channel := channelmodel.ChannelInstagram{ AccountID: 1, InboxID: inbox.ID, InstagramAccountID: "ig-123", InstagramBusinessAccountID: "ig-business-123", PageAccessToken: "page-token", ConnectedFBPageID: "page-123", InstagramAccountName: "gochat", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create instagram channel: %v", err) } t.Setenv("INSTAGRAM_APP_SECRET", "ig-secret") body := []byte(`{"object":"instagram","entry":[{"id":"ig-123","time":1,"messaging":[{"sender":{"id":"user-1"},"recipient":{"id":"ig-123"},"timestamp":1,"message":{"mid":"mid-1","text":"hello"}}]}]}`) h := NewFacebookWebhookHandler(nil, nil, db) r := gin.New() r.POST("/webhooks/instagram", h.HandleInstagramWebhook) req := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) req.Header.Set("X-Hub-Signature-256", metaSignature("ig-secret", body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusOK { t.Fatalf("expected 200, got %d body=%s", w.Code, w.Body.String()) } assertPersistedMessage(t, db, inbox.ID, "mid-1", "hello") } func TestInstagramWebhookRejectsMissingSignature(t *testing.T) { gin.SetMode(gin.TestMode) db := newWebhookLookupTestDB(t) inbox := seedWebhookInbox(t, db, "instagram") channel := channelmodel.ChannelInstagram{ AccountID: 1, InboxID: inbox.ID, InstagramAccountID: "ig-123", InstagramBusinessAccountID: "ig-business-123", PageAccessToken: "page-token", ConnectedFBPageID: "page-123", InstagramAccountName: "gochat", } if err := db.Create(&channel).Error; err != nil { t.Fatalf("create instagram channel: %v", err) } t.Setenv("INSTAGRAM_APP_SECRET", "ig-secret") body := []byte(`{"object":"instagram","entry":[{"id":"ig-123","time":1,"messaging":[{"sender":{"id":"user-1"},"recipient":{"id":"ig-123"},"timestamp":1,"message":{"mid":"mid-missing-sig","text":"hello"}}]}]}`) h := NewFacebookWebhookHandler(nil, nil, db) r := gin.New() r.POST("/webhooks/instagram", h.HandleInstagramWebhook) req := httptest.NewRequest(http.MethodPost, "/webhooks/instagram", bytes.NewReader(body)) w := httptest.NewRecorder() r.ServeHTTP(w, req) if w.Code != http.StatusUnauthorized { t.Fatalf("expected 401, got %d body=%s", w.Code, w.Body.String()) } var count int64 if err := db.Model(&model.Message{}).Where("inbox_id = ? AND source_id = ?", inbox.ID, "mid-missing-sig").Count(&count).Error; err != nil { t.Fatalf("count message: %v", err) } if count != 0 { t.Fatalf("expected no persisted message, got %d", count) } } func assertPersistedMessage(t *testing.T, db *gorm.DB, inboxID uint, sourceID string, content string) model.Message { t.Helper() var message model.Message if err := db.Where("inbox_id = ? AND source_id = ?", inboxID, sourceID).First(&message).Error; err != nil { t.Fatalf("expected message source_id=%s persisted: %v", sourceID, err) } if message.Content != content || message.MessageType != string(model.MessageTypeIncoming) { t.Fatalf("unexpected message source_id=%s: %#v", sourceID, message) } return message }