## 核心修复 ### 1. Auto-Reply Sender 修复(所有渠道) - AutoReplyListener.sendAutoReply() 通过 botInboxRepo 查询 inbox 关联的 AgentBot - 使用正确的 SenderType="AgentBot"(非小写 agent_bot)传递真实 AgentBot ID - bootstrap 注入 agentBotInboxRepo/agentBotRepo 依赖 ### 2. 事件数据 BUG 修复(影响所有 Webhook 渠道) - incoming_persister.dispatch(): 补全 sender_type/content 到 event.Data - channel/webhook.go: HandleWebhook 同步分发也补全 sender_type/content - 未补全前 AutoReplyListener 找不到字段直接跳过 ### 3. Web Widget SDK 生产验证修复 - cookie → localStorage token 同步(frontend/index.html) - 路由双注册修复(router.go) - Vite SPA 模式 + /widget 重写(vite.config.ts) - WidgetService 注入 Dispatcher 触发事件分发 ### 4. LLM 真实模型对接 - 配置 deepseek-v4-flash @ http://10.58.144.6:2014/v1 - LLM-mode auto-reply 规则创建并验证通过 - Prompt 文档落地: docs/captain-ai-auto-replay-prompt.md ### 5. 新增基础设施 - Helm chart (deploy/helm/) - Widget SDK 生产测试页面 - QA 报告 Closes: BUG-W2 (auth sync), BUG-W3 (route double-reg), BUG-WEBHOOK-EVENT (missing event data fields)
632 lines
21 KiB
Go
632 lines
21 KiB
Go
package webhook
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"gorm.io/datatypes"
|
|
"gorm.io/gorm"
|
|
|
|
"github.com/gochat/gochat/internal/channel"
|
|
"github.com/gochat/gochat/internal/model"
|
|
"github.com/gochat/gochat/internal/worker"
|
|
applogger "github.com/gochat/gochat/pkg/logger"
|
|
)
|
|
|
|
// IncomingPersister is the durable boundary after provider-specific webhook parsing.
|
|
// Reference: Chatwoot IncomingMessageService creates ContactInbox, Conversation, and Message.
|
|
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 {
|
|
Contact *model.Contact
|
|
ContactInbox *model.ContactInbox
|
|
Conversation *model.Conversation
|
|
Message *model.Message
|
|
Duplicate bool
|
|
ContactCreated bool
|
|
ConversationCreated bool
|
|
}
|
|
|
|
func (p *IncomingPersister) SetWorkerPool(wp *worker.WorkerPool) *IncomingPersister {
|
|
if p == nil || wp == nil {
|
|
return p
|
|
}
|
|
p.worker = wp
|
|
p.registerJobs(wp)
|
|
return p
|
|
}
|
|
|
|
func (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
|
|
}
|
|
p := &IncomingPersister{db: db}
|
|
if len(dispatcher) > 0 {
|
|
p.dispatcher = dispatcher[0]
|
|
}
|
|
return p
|
|
}
|
|
|
|
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 err := validateIncomingMessage(p, inbox, msg); err != nil {
|
|
return nil, err
|
|
}
|
|
if p.worker != nil {
|
|
return nil, p.enqueueIncomingMessagePersist(ctx, inbox.ID, msg)
|
|
}
|
|
return p.performPersistIncoming(ctx, inbox, msg)
|
|
}
|
|
|
|
func validateIncomingMessage(p *IncomingPersister, inbox *model.Inbox, msg *channel.IncomingMessage) error {
|
|
if p == nil || p.db == nil || inbox == nil || msg == nil {
|
|
return nil
|
|
}
|
|
if msg.SourceID == "" {
|
|
return fmt.Errorf("incoming message missing source_id")
|
|
}
|
|
senderID := msg.SenderID
|
|
if senderID == "" {
|
|
senderID = msg.ConversationID
|
|
}
|
|
if senderID == "" {
|
|
return fmt.Errorf("incoming message missing sender_id")
|
|
}
|
|
if msg.Content == "" && len(msg.Attachments) == 0 {
|
|
return fmt.Errorf("incoming message has no content or attachments")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *IncomingPersister) performPersistIncoming(ctx context.Context, inbox *model.Inbox, msg *channel.IncomingMessage) (*IncomingPersistResult, error) {
|
|
if p == nil || p.db == nil || inbox == nil || msg == nil {
|
|
return nil, nil
|
|
}
|
|
senderID := msg.SenderID
|
|
if senderID == "" {
|
|
senderID = msg.ConversationID
|
|
}
|
|
|
|
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, contactCreated, err := p.resolveOrCreateContactInbox(ctx, tx, inbox, msg, senderID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
result.Contact = contact
|
|
result.ContactInbox = contactInbox
|
|
result.ContactCreated = contactCreated
|
|
|
|
conversation, conversationCreated, err := p.resolveOrCreateConversation(ctx, tx, inbox, contact, contactInbox, msg)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
result.Conversation = conversation
|
|
result.ConversationCreated = conversationCreated
|
|
|
|
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
|
|
}
|
|
p.indexIncomingResult(ctx, &result)
|
|
p.dispatchIncomingEvents(ctx, inbox, &result)
|
|
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 {
|
|
return p.UpdateMessageStatusWithError(ctx, inbox, sourceID, status, occurredAt, "")
|
|
}
|
|
|
|
// UpdateMessageStatusWithError applies provider receipts and records provider failure details.
|
|
func (p *IncomingPersister) UpdateMessageStatusWithError(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time, externalError string) error {
|
|
if p == nil || p.db == nil || inbox == nil || sourceID == "" {
|
|
return nil
|
|
}
|
|
if p.worker != nil {
|
|
return p.enqueueMessageStatusUpdate(ctx, inbox.ID, sourceID, status, occurredAt, externalError)
|
|
}
|
|
return p.performMessageStatusUpdate(ctx, inbox, sourceID, status, occurredAt, externalError)
|
|
}
|
|
|
|
func (p *IncomingPersister) performMessageStatusUpdate(ctx context.Context, inbox *model.Inbox, sourceID string, status model.MessageStatus, occurredAt *time.Time, externalError string) error {
|
|
if p == nil || p.db == nil || inbox == nil || sourceID == "" {
|
|
return nil
|
|
}
|
|
|
|
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 {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
if !validProviderMessageStatusTransition(model.MessageStatus(message.Status), status) {
|
|
return nil
|
|
}
|
|
updates := map[string]any{
|
|
"status": string(status),
|
|
"content_attributes": setProviderMessageExternalError(message.ContentAttributes, status, externalError),
|
|
}
|
|
if err := tx.Model(&message).Updates(updates).Error; err != nil {
|
|
return err
|
|
}
|
|
message.Status = string(status)
|
|
message.ContentAttributes = updates["content_attributes"].(datatypes.JSON)
|
|
contactID, err := p.deliveryStatusContactID(ctx, tx, &message)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if contactID == 0 {
|
|
p.dispatchMessageStatusEvent(ctx, inbox, &message)
|
|
copy := message
|
|
updatedMessage = ©
|
|
return nil
|
|
}
|
|
if err := p.upsertDeliveryStatus(ctx, tx, &message, contactID, status, occurredAt); err != nil {
|
|
return err
|
|
}
|
|
p.dispatchMessageStatusEvent(ctx, inbox, &message)
|
|
copy := message
|
|
updatedMessage = ©
|
|
return nil
|
|
}); err != nil {
|
|
return err
|
|
}
|
|
p.indexMessage(ctx, updatedMessage)
|
|
return nil
|
|
}
|
|
|
|
// 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
|
|
}
|
|
if p.worker != nil {
|
|
return p.enqueueContactMessagesStatusUpdate(ctx, inbox.ID, contactSourceID, status, occurredAt)
|
|
}
|
|
return p.performContactMessagesStatusUpdate(ctx, inbox, contactSourceID, status, occurredAt)
|
|
}
|
|
|
|
func (p *IncomingPersister) performContactMessagesStatusUpdate(ctx context.Context, inbox *model.Inbox, contactSourceID string, status model.MessageStatus, occurredAt *time.Time) error {
|
|
if p == nil || p.db == nil || inbox == nil || contactSourceID == "" {
|
|
return nil
|
|
}
|
|
|
|
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 {
|
|
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)
|
|
}
|
|
|
|
var messages []model.Message
|
|
if err := query.Select("messages.*").Find(&messages).Error; err != nil {
|
|
return err
|
|
}
|
|
for i := range messages {
|
|
currentStatus := model.MessageStatus(messages[i].Status)
|
|
if currentStatus == status {
|
|
if err := p.upsertDeliveryStatus(ctx, tx, &messages[i], contactInbox.ContactID, status, occurredAt); err != nil {
|
|
return err
|
|
}
|
|
updatedMessages = append(updatedMessages, messages[i])
|
|
continue
|
|
}
|
|
if !validProviderMessageStatusTransition(currentStatus, status) {
|
|
continue
|
|
}
|
|
if err := tx.Model(&messages[i]).Updates(map[string]any{
|
|
"status": string(status),
|
|
"content_attributes": setProviderMessageExternalError(messages[i].ContentAttributes, status, ""),
|
|
}).Error; err != nil {
|
|
return err
|
|
}
|
|
messages[i].Status = string(status)
|
|
if err := p.upsertDeliveryStatus(ctx, tx, &messages[i], contactInbox.ContactID, status, occurredAt); err != nil {
|
|
return err
|
|
}
|
|
p.dispatchMessageStatusEvent(ctx, inbox, &messages[i])
|
|
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) {
|
|
if message == nil {
|
|
return 0, nil
|
|
}
|
|
if message.ConversationID != 0 {
|
|
var conversation model.Conversation
|
|
if err := tx.WithContext(ctx).Select("id", "contact_id").First(&conversation, message.ConversationID).Error; err != nil {
|
|
if err != gorm.ErrRecordNotFound {
|
|
return 0, err
|
|
}
|
|
} else if conversation.ContactID != 0 {
|
|
return conversation.ContactID, nil
|
|
}
|
|
}
|
|
if message.SenderID != nil {
|
|
return *message.SenderID, nil
|
|
}
|
|
return 0, nil
|
|
}
|
|
|
|
func (p *IncomingPersister) upsertDeliveryStatus(ctx context.Context, tx *gorm.DB, message *model.Message, contactID uint, status model.MessageStatus, occurredAt *time.Time) error {
|
|
var delivery model.DeliveryStatus
|
|
err := tx.WithContext(ctx).Where("message_id = ? AND contact_id = ?", message.ID, contactID).First(&delivery).Error
|
|
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) dispatchIncomingEvents(ctx context.Context, inbox *model.Inbox, result *IncomingPersistResult) {
|
|
if p.dispatcher == nil || result == nil || result.Duplicate {
|
|
return
|
|
}
|
|
if result.ContactCreated && result.Contact != nil {
|
|
p.dispatch(ctx, channel.EventContactCreated, inbox, result, nil)
|
|
}
|
|
if result.Conversation != nil {
|
|
if result.ConversationCreated {
|
|
p.dispatch(ctx, channel.EventConversationCreated, inbox, result, nil)
|
|
p.dispatch(ctx, channel.EventConversationOpened, inbox, result, nil)
|
|
} else {
|
|
p.dispatch(ctx, channel.EventConversationUpdated, inbox, result, nil)
|
|
}
|
|
}
|
|
if result.Message != nil {
|
|
p.dispatch(ctx, channel.EventMessageCreated, inbox, result, result.Message)
|
|
p.dispatch(ctx, channel.EventMessageIncoming, inbox, result, result.Message)
|
|
}
|
|
}
|
|
|
|
func (p *IncomingPersister) dispatchMessageStatusEvent(ctx context.Context, inbox *model.Inbox, message *model.Message) {
|
|
if p.dispatcher == nil || inbox == nil || message == nil {
|
|
return
|
|
}
|
|
result := &IncomingPersistResult{Message: message}
|
|
p.dispatch(ctx, channel.EventMessageStatusUpdated, inbox, result, message)
|
|
}
|
|
|
|
func (p *IncomingPersister) dispatch(ctx context.Context, eventType channel.EventType, inbox *model.Inbox, result *IncomingPersistResult, message *model.Message) {
|
|
if p.dispatcher == nil || inbox == nil {
|
|
return
|
|
}
|
|
channelType := channel.ChannelType(inbox.ChannelType)
|
|
event := channel.NewChannelEvent(eventType, channelType, inbox.AccountID, inbox.ID)
|
|
if result != nil {
|
|
if result.Conversation != nil {
|
|
event.ConversationID = result.Conversation.ID
|
|
event.Data["conversation"] = result.Conversation
|
|
}
|
|
if result.Contact != nil {
|
|
event.ContactID = result.Contact.ID
|
|
event.Data["contact"] = result.Contact
|
|
}
|
|
if result.Message != nil {
|
|
event.Data["message"] = result.Message
|
|
event.Data["sender_type"] = result.Message.SenderType
|
|
event.Data["content"] = result.Message.Content
|
|
}
|
|
}
|
|
if message != nil {
|
|
event.ConversationID = message.ConversationID
|
|
event.Data["message"] = message
|
|
event.Data["sender_type"] = message.SenderType
|
|
event.Data["content"] = message.Content
|
|
if message.SenderID != nil {
|
|
event.ContactID = *message.SenderID
|
|
}
|
|
}
|
|
event.Data["inbox"] = inbox
|
|
if err := p.dispatcher.DispatchAsync(ctx, event); err != nil {
|
|
// Chatwoot's async side effects should not make provider webhooks fail.
|
|
_ = err
|
|
applogger.L().Warnf("IncomingPersister: failed to dispatch %s event for inbox=%d: %v",
|
|
eventType, inbox.ID, err)
|
|
}
|
|
}
|
|
|
|
func (p *IncomingPersister) resolveOrCreateContactInbox(ctx context.Context, tx *gorm.DB, inbox *model.Inbox, msg *channel.IncomingMessage, senderID string) (*model.Contact, *model.ContactInbox, bool, 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, false, err
|
|
}
|
|
if err := tx.WithContext(ctx).First(&contact, contact.ID).Error; err != nil {
|
|
return nil, nil, false, err
|
|
}
|
|
}
|
|
return &contact, &contactInbox, false, nil
|
|
} else if err != gorm.ErrRecordNotFound {
|
|
return nil, nil, false, 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, false, 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, false, err
|
|
}
|
|
return &contact, &contactInbox, true, 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, bool, 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, false, nil
|
|
} else if err != gorm.ErrRecordNotFound {
|
|
return nil, false, 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, false, err
|
|
}
|
|
return &conversation, true, 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 replyToExternalID := strings.TrimSpace(msg.ReplyToID); replyToExternalID != "" {
|
|
// Chatwoot stores provider reply references as in_reply_to_external_id,
|
|
// then resolves the same-conversation message by source_id before save.
|
|
// Keep in_reply_to strictly as the internal numeric message ID expected by
|
|
// the dashboard and its cursor-based message API.
|
|
var replyMessage model.Message
|
|
err := tx.WithContext(ctx).
|
|
Where("conversation_id = ? AND source_id = ?", conversation.ID, replyToExternalID).
|
|
First(&replyMessage).Error
|
|
if err == nil {
|
|
contentAttrs["in_reply_to"] = replyMessage.ID
|
|
contentAttrs["in_reply_to_external_id"] = replyMessage.SourceID
|
|
} else if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
// Match Chatwoot's InReplyToMessageBuilder: unresolved references
|
|
// are normalized to null instead of being persisted as an invalid
|
|
// internal message ID.
|
|
contentAttrs["in_reply_to"] = nil
|
|
contentAttrs["in_reply_to_external_id"] = nil
|
|
} else {
|
|
return nil, err
|
|
}
|
|
}
|
|
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)
|
|
}
|