feat(search): index provider webhooks

This commit is contained in:
2026-06-07 15:23:29 +08:00
parent 28e7205410
commit 8eabadc3fc
10 changed files with 183 additions and 13 deletions
@@ -76,6 +76,13 @@ func (h *FacebookWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *Facebook
return h
}
func (h *FacebookWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *FacebookWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
// HandleFacebookVerification handles the GET webhook verification request from Facebook.
// URL pattern: /webhooks/facebook/:inbox_id
// Method: GET
+80 -4
View File
@@ -21,6 +21,13 @@ type IncomingPersister struct {
db *gorm.DB
dispatcher *channel.Dispatcher
worker *worker.WorkerPool
indexer IncomingSearchIndexer
}
type IncomingSearchIndexer interface {
IndexConversation(ctx context.Context, conversation *model.Conversation) error
IndexMessage(ctx context.Context, message *model.Message) error
IndexContact(ctx context.Context, contact *model.Contact) error
}
type IncomingPersistResult struct {
@@ -42,6 +49,14 @@ func (p *IncomingPersister) SetWorkerPool(wp *worker.WorkerPool) *IncomingPersis
return p
}
func (p *IncomingPersister) SetSearchIndexer(indexer IncomingSearchIndexer) *IncomingPersister {
if p == nil {
return p
}
p.indexer = indexer
return p
}
func NewIncomingPersister(db *gorm.DB, dispatcher ...*channel.Dispatcher) *IncomingPersister {
if db == nil {
return nil
@@ -131,6 +146,7 @@ func (p *IncomingPersister) performPersistIncoming(ctx context.Context, inbox *m
if err != nil {
return nil, err
}
p.indexIncomingResult(ctx, &result)
p.dispatchIncomingEvents(ctx, inbox, &result)
return &result, nil
}
@@ -156,7 +172,8 @@ func (p *IncomingPersister) performMessageStatusUpdate(ctx context.Context, inbo
return nil
}
return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var updatedMessage *model.Message
if err := p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var message model.Message
if err := tx.Where("inbox_id = ? AND source_id = ?", inbox.ID, sourceID).First(&message).Error; err != nil {
if err == gorm.ErrRecordNotFound {
@@ -182,14 +199,22 @@ func (p *IncomingPersister) performMessageStatusUpdate(ctx context.Context, inbo
}
if contactID == 0 {
p.dispatchMessageStatusEvent(ctx, inbox, &message)
copy := message
updatedMessage = &copy
return nil
}
if err := p.upsertDeliveryStatus(ctx, tx, &message, contactID, status, occurredAt); err != nil {
return err
}
p.dispatchMessageStatusEvent(ctx, inbox, &message)
copy := message
updatedMessage = &copy
return nil
})
}); err != nil {
return err
}
p.indexMessage(ctx, updatedMessage)
return nil
}
// UpdateContactConversationMessagesStatus applies read receipts that only identify the contact/conversation.
@@ -208,7 +233,8 @@ func (p *IncomingPersister) performContactMessagesStatusUpdate(ctx context.Conte
return nil
}
return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
updatedMessages := []model.Message{}
if err := p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var contactInbox model.ContactInbox
if err := tx.Where("inbox_id = ? AND source_id = ?", inbox.ID, contactSourceID).First(&contactInbox).Error; err != nil {
if err == gorm.ErrRecordNotFound {
@@ -244,9 +270,59 @@ func (p *IncomingPersister) performContactMessagesStatusUpdate(ctx context.Conte
return err
}
p.dispatchMessageStatusEvent(ctx, inbox, &messages[i])
updatedMessages = append(updatedMessages, messages[i])
}
return nil
})
}); err != nil {
return err
}
for i := range updatedMessages {
p.indexMessage(ctx, &updatedMessages[i])
}
return nil
}
func (p *IncomingPersister) indexIncomingResult(ctx context.Context, result *IncomingPersistResult) {
if p == nil || p.indexer == nil || result == nil || result.Duplicate {
return
}
if result.Contact != nil {
p.indexContact(ctx, result.Contact)
}
if result.Conversation != nil {
p.indexConversation(ctx, result.Conversation)
}
if result.Message != nil {
p.indexMessage(ctx, result.Message)
}
}
func (p *IncomingPersister) indexContact(ctx context.Context, contact *model.Contact) {
if p == nil || p.indexer == nil || contact == nil {
return
}
if err := p.indexer.IndexContact(ctx, contact); err != nil {
// Provider webhooks should not fail because search is temporarily unavailable.
_ = err
}
}
func (p *IncomingPersister) indexConversation(ctx context.Context, conversation *model.Conversation) {
if p == nil || p.indexer == nil || conversation == nil {
return
}
if err := p.indexer.IndexConversation(ctx, conversation); err != nil {
_ = err
}
}
func (p *IncomingPersister) indexMessage(ctx context.Context, message *model.Message) {
if p == nil || p.indexer == nil || message == nil {
return
}
if err := p.indexer.IndexMessage(ctx, message); err != nil {
_ = err
}
}
func (p *IncomingPersister) deliveryStatusContactID(ctx context.Context, tx *gorm.DB, message *model.Message) (uint, error) {
+7
View File
@@ -50,6 +50,13 @@ func (h *LineWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *LineWebhookH
return h
}
func (h *LineWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *LineWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
// HandleLineWebhook processes an incoming LINE webhook Gin request.
func (h *LineWebhookHandler) HandleLineWebhook(c *gin.Context) {
lineChannelID := c.Param("line_channel_id")
@@ -63,6 +63,13 @@ func (h *TelegramWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *Telegram
return h
}
func (h *TelegramWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *TelegramWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
// HandleTelegramWebhook processes an incoming Telegram webhook Gin request.
// URL pattern: /webhooks/telegram/:bot_token
// Method: POST
@@ -54,6 +54,13 @@ func (h *TikTokWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *TikTokWebh
return h
}
func (h *TikTokWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *TikTokWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
// HandleTikTokWebhook processes incoming TikTok webhook HTTP requests.
func (h *TikTokWebhookHandler) HandleTikTokWebhook(c *gin.Context) {
body, err := io.ReadAll(c.Request.Body)
@@ -41,6 +41,13 @@ func (h *TwilioWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *TwilioWebh
return h
}
func (h *TwilioWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *TwilioWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
// NewTwilioWebhookHandler creates a Twilio SMS webhook handler for Gin integration.
func NewTwilioWebhookHandler(twilioWebhook *twiliochannel.WebhookHandler, db *gorm.DB, dispatcher ...*channel.Dispatcher) *TwilioWebhookHandler {
return &TwilioWebhookHandler{
@@ -35,6 +35,27 @@ 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
}
@@ -175,8 +196,9 @@ func TestIncomingPersisterQueuesMessageStatusUpdateWithWorker(t *testing.T) {
dispatcher := channel.NewDispatcher()
listener := &recordingListener{}
dispatcher.Register(listener)
indexer := &recordingIncomingSearchIndexer{}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db, dispatcher)
persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer)
msg := &channel.IncomingMessage{
ChannelType: channel.ChannelTelegram,
@@ -232,6 +254,9 @@ func TestIncomingPersisterQueuesMessageStatusUpdateWithWorker(t *testing.T) {
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) {
@@ -250,7 +275,8 @@ func TestIncomingPersisterStatusJobDoesNotDowngradeRead(t *testing.T) {
t.Fatalf("create message: %v", err)
}
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db).SetWorkerPool(wp)
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)
@@ -330,7 +356,8 @@ func TestIncomingPersisterQueuesContactMessagesStatusUpdateWithWorker(t *testing
alreadyRead := createWebhookStatusMessage(t, db, inbox, conversation, "read-before", model.MessageStatusRead, cutoff.Add(-20*time.Second))
afterCutoff := createWebhookStatusMessage(t, db, inbox, conversation, "sent-after", model.MessageStatusSent, cutoff.Add(time.Minute))
wp := worker.NewWorkerPoolWithOptions(db, worker.WithQueues("low"))
persister := NewIncomingPersister(db).SetWorkerPool(wp)
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)
@@ -355,6 +382,9 @@ func TestIncomingPersisterQueuesContactMessagesStatusUpdateWithWorker(t *testing
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) {
@@ -363,7 +393,8 @@ func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) {
dispatcher := channel.NewDispatcher()
listener := &recordingListener{}
dispatcher.Register(listener)
persister := NewIncomingPersister(db, dispatcher)
indexer := &recordingIncomingSearchIndexer{}
persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer)
msg := &channel.IncomingMessage{
ChannelType: channel.ChannelTelegram,
@@ -387,6 +418,9 @@ func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) {
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 {
@@ -395,6 +429,9 @@ func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) {
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"
@@ -405,6 +442,9 @@ func TestIncomingPersisterCreatesConversationMessageAndDedupes(t *testing.T) {
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 {
@@ -426,8 +466,9 @@ func TestIncomingPersisterQueuesIncomingMessageWithWorker(t *testing.T) {
dispatcher := channel.NewDispatcher()
listener := &recordingListener{}
dispatcher.Register(listener)
indexer := &recordingIncomingSearchIndexer{}
wp := worker.NewWorkerPool(db)
persister := NewIncomingPersister(db, dispatcher).SetWorkerPool(wp)
persister := NewIncomingPersister(db, dispatcher).SetSearchIndexer(indexer).SetWorkerPool(wp)
msg := &channel.IncomingMessage{
ChannelType: channel.ChannelTelegram,
@@ -464,6 +505,9 @@ func TestIncomingPersisterQueuesIncomingMessageWithWorker(t *testing.T) {
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)
@@ -78,6 +78,13 @@ func (h *WhatsAppWebhookHandler) WithWorkerPool(wp *worker.WorkerPool) *WhatsApp
return h
}
func (h *WhatsAppWebhookHandler) WithSearchIndexer(indexer IncomingSearchIndexer) *WhatsAppWebhookHandler {
if h != nil && h.persister != nil {
h.persister.SetSearchIndexer(indexer)
}
return h
}
type whatsAppPersisterAdapter struct {
persister *IncomingPersister
}