Files
gochat/internal/handler/webhook/incoming_persister.go
T

329 lines
11 KiB
Go

package webhook
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/google/uuid"
"gorm.io/datatypes"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
)
// IncomingPersister is the durable boundary after provider-specific webhook parsing.
// Reference: Chatwoot IncomingMessageService creates ContactInbox, Conversation, and Message.
type IncomingPersister struct {
db *gorm.DB
}
type IncomingPersistResult struct {
Contact *model.Contact
ContactInbox *model.ContactInbox
Conversation *model.Conversation
Message *model.Message
Duplicate bool
}
func NewIncomingPersister(db *gorm.DB) *IncomingPersister {
if db == nil {
return nil
}
return &IncomingPersister{db: db}
}
func (p *IncomingPersister) PersistIncoming(ctx context.Context, inbox *model.Inbox, msg *channel.IncomingMessage) (*IncomingPersistResult, error) {
if p == nil || p.db == nil || inbox == nil || msg == nil {
return nil, nil
}
if msg.SourceID == "" {
return nil, fmt.Errorf("incoming message missing source_id")
}
senderID := msg.SenderID
if senderID == "" {
senderID = msg.ConversationID
}
if senderID == "" {
return nil, fmt.Errorf("incoming message missing sender_id")
}
if msg.Content == "" && len(msg.Attachments) == 0 {
return nil, fmt.Errorf("incoming message has no content or attachments")
}
var result IncomingPersistResult
err := p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var existingMessage model.Message
if err := tx.Where("inbox_id = ? AND source_id = ?", inbox.ID, msg.SourceID).First(&existingMessage).Error; err == nil {
result.Message = &existingMessage
result.Duplicate = true
return nil
} else if err != gorm.ErrRecordNotFound {
return err
}
contact, contactInbox, err := p.resolveOrCreateContactInbox(ctx, tx, inbox, msg, senderID)
if err != nil {
return err
}
result.Contact = contact
result.ContactInbox = contactInbox
conversation, err := p.resolveOrCreateConversation(ctx, tx, inbox, contact, contactInbox, msg)
if err != nil {
return err
}
result.Conversation = conversation
message, err := p.createMessage(ctx, tx, inbox, conversation, contact, msg)
if err != nil {
return err
}
result.Message = message
return nil
})
if err != nil {
return nil, err
}
return &result, nil
}
// UpdateMessageStatus applies provider delivery/read/failed receipts to existing messages.
func (p *IncomingPersister) UpdateMessageStatus(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time) error {
if p == nil || p.db == nil || inbox == nil || sourceID == "" {
return nil
}
return 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 {
return nil
}
return err
}
if err := tx.Model(&message).Update("status", string(status)).Error; err != nil {
return err
}
if message.SenderID == nil {
return nil
}
return p.upsertDeliveryStatus(ctx, tx, &message, *message.SenderID, status, occurredAt)
})
}
// UpdateContactConversationMessagesStatus applies read receipts that only identify the contact/conversation.
func (p *IncomingPersister) UpdateContactConversationMessagesStatus(ctx context.Context, inbox *model.Inbox, contactSourceID string, status model.MessageStatus, occurredAt *time.Time) error {
if p == nil || p.db == nil || inbox == nil || contactSourceID == "" {
return nil
}
return p.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var contactInbox model.ContactInbox
if err := tx.Where("inbox_id = ? AND source_id = ?", inbox.ID, contactSourceID).First(&contactInbox).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return nil
}
return err
}
query := tx.Model(&model.Message{}).
Joins("JOIN conversations ON conversations.id = messages.conversation_id").
Where("messages.inbox_id = ? AND conversations.contact_id = ?", inbox.ID, contactInbox.ContactID).
Where("messages.message_type = ?", model.MessageTypeOutgoing)
if occurredAt != nil {
query = query.Where("messages.created_at <= ?", *occurredAt)
}
return query.Update("status", string(status)).Error
})
}
func (p *IncomingPersister) upsertDeliveryStatus(ctx context.Context, tx *gorm.DB, message *model.Message, contactID uint, status model.MessageStatus, occurredAt *time.Time) error {
var delivery model.DeliveryStatus
err := tx.WithContext(ctx).Where("message_id = ? AND contact_id = ?", message.ID, contactID).First(&delivery).Error
if err != nil && err != gorm.ErrRecordNotFound {
return err
}
if err == gorm.ErrRecordNotFound {
delivery = model.DeliveryStatus{MessageID: message.ID, InboxID: message.InboxID, ContactID: contactID}
}
delivery.Status = status
if occurredAt != nil {
switch status {
case model.MessageStatusRead:
delivery.ReadAt = occurredAt
case model.MessageStatusDelivered:
delivery.DeliveredAt = occurredAt
}
}
if delivery.ID == 0 {
return tx.WithContext(ctx).Create(&delivery).Error
}
return tx.WithContext(ctx).Save(&delivery).Error
}
func (p *IncomingPersister) resolveOrCreateContactInbox(ctx context.Context, tx *gorm.DB, inbox *model.Inbox, msg *channel.IncomingMessage, senderID string) (*model.Contact, *model.ContactInbox, error) {
var contactInbox model.ContactInbox
if err := tx.WithContext(ctx).Preload("Contact").Where("inbox_id = ? AND source_id = ?", inbox.ID, senderID).First(&contactInbox).Error; err == nil {
contact := contactInbox.Contact
updates := map[string]interface{}{}
if msg.SenderName != "" && contact.Name != msg.SenderName {
updates["name"] = msg.SenderName
}
if len(updates) > 0 {
if err := tx.WithContext(ctx).Model(&contact).Updates(updates).Error; err != nil {
return nil, nil, err
}
if err := tx.WithContext(ctx).First(&contact, contact.ID).Error; err != nil {
return nil, nil, err
}
}
return &contact, &contactInbox, nil
} else if err != gorm.ErrRecordNotFound {
return nil, nil, err
}
name := msg.SenderName
if name == "" {
name = senderID
}
attrs := mergeChannelConfig(msg.SenderExtra, map[string]interface{}{
"source_id": senderID,
"channel_type": string(msg.ChannelType),
})
contact := model.Contact{
AccountID: inbox.AccountID,
Name: name,
Identifier: senderID,
SourceID: string(msg.ChannelType),
AdditionalAttributes: mustJSON(attrs),
CustomAttributes: datatypes.JSON("{}"),
}
if err := tx.WithContext(ctx).Create(&contact).Error; err != nil {
return nil, nil, err
}
contactInbox = model.ContactInbox{
ContactID: contact.ID,
InboxID: inbox.ID,
SourceID: senderID,
PubsubToken: uuid.NewString(),
}
if err := tx.WithContext(ctx).Create(&contactInbox).Error; err != nil {
return nil, nil, err
}
return &contact, &contactInbox, nil
}
func (p *IncomingPersister) resolveOrCreateConversation(ctx context.Context, tx *gorm.DB, inbox *model.Inbox, contact *model.Contact, contactInbox *model.ContactInbox, msg *channel.IncomingMessage) (*model.Conversation, error) {
var conversation model.Conversation
if err := tx.WithContext(ctx).
Where("account_id = ? AND inbox_id = ? AND contact_id = ? AND status = ?", inbox.AccountID, inbox.ID, contact.ID, model.ConversationStatusOpen).
Order("id DESC").First(&conversation).Error; err == nil {
now := time.Now().Unix()
_ = tx.WithContext(ctx).Model(&conversation).Updates(map[string]interface{}{
"last_activity_at": now,
"last_message_at": now,
}).Error
return &conversation, nil
} else if err != gorm.ErrRecordNotFound {
return nil, err
}
now := time.Now().Unix()
channelType := inbox.ChannelType
if msg.ChannelType != "" {
channelType = string(msg.ChannelType)
}
conversation = model.Conversation{
AccountID: inbox.AccountID,
InboxID: inbox.ID,
ContactID: contact.ID,
ContactInboxID: &contactInbox.ID,
Status: string(model.ConversationStatusOpen),
Priority: "none",
ChannelType: channelType,
Channel: inbox.ChannelType,
LastActivityAt: &now,
LastMessageAt: &now,
LastNonSysMsgAt: &now,
AdditionalAttributes: mustJSON(mergeChannelConfig(msg.ConversationExtra, map[string]interface{}{"external_conversation_id": msg.ConversationID})),
CustomAttributes: datatypes.JSON("{}"),
}
if err := tx.WithContext(ctx).Create(&conversation).Error; err != nil {
return nil, err
}
return &conversation, nil
}
func (p *IncomingPersister) createMessage(ctx context.Context, tx *gorm.DB, inbox *model.Inbox, conversation *model.Conversation, contact *model.Contact, msg *channel.IncomingMessage) (*model.Message, error) {
contentAttrs := map[string]interface{}{}
if len(msg.Attachments) > 0 {
contentAttrs["attachments"] = msg.Attachments
}
if msg.ReplyToID != "" {
contentAttrs["in_reply_to"] = msg.ReplyToID
}
message := model.Message{
ConversationID: conversation.ID,
AccountID: inbox.AccountID,
InboxID: inbox.ID,
SenderID: &contact.ID,
SenderType: string(model.SenderTypeContact),
Content: msg.Content,
ContentType: mapIncomingContentType(msg.ContentType),
Status: string(model.MessageStatusSent),
MessageType: string(model.MessageTypeIncoming),
SourceID: msg.SourceID,
ContentAttributes: mustJSON(contentAttrs),
AdditionalAttributes: mustJSON(msg.Extra),
}
if err := tx.WithContext(ctx).Create(&message).Error; err != nil {
return nil, err
}
return &message, nil
}
func mapIncomingContentType(contentType channel.ContentType) string {
switch contentType {
case channel.ContentImage:
return string(model.MessageContentTypeImage)
case channel.ContentAudio:
return string(model.MessageContentTypeAudio)
case channel.ContentVideo:
return string(model.MessageContentTypeVideo)
case channel.ContentFile:
return string(model.MessageContentTypeFile)
case channel.ContentLocation:
return string(model.MessageContentTypeLocation)
default:
return string(model.MessageContentTypeText)
}
}
func mergeChannelConfig(base channel.ChannelConfig, extra map[string]interface{}) map[string]interface{} {
merged := map[string]interface{}{}
for k, v := range base {
merged[k] = v
}
for k, v := range extra {
if v != "" && v != nil {
merged[k] = v
}
}
return merged
}
func mustJSON(value interface{}) datatypes.JSON {
if value == nil {
return datatypes.JSON("{}")
}
data, err := json.Marshal(value)
if err != nil || len(data) == 0 {
return datatypes.JSON("{}")
}
return datatypes.JSON(data)
}