feat(realtime): complete fake channel reply flow

This commit is contained in:
2026-07-13 14:57:28 +08:00
parent a712b91982
commit d16bb55d91
24 changed files with 1441 additions and 840 deletions
+29 -17
View File
@@ -110,16 +110,16 @@ func (p *FakeProvider) OnDestroy(ctx context.Context, inbox *model.Inbox, config
// FakeIncomingPayload is the JSON body FakeMessagePlatform posts to
// /webhooks/fake/:identifier.
type FakeIncomingPayload struct {
Event string `json:"event"`
MessageID string `json:"message_id"`
SenderID string `json:"sender_id"`
SenderName string `json:"sender_name"`
Content string `json:"content"`
ContentType string `json:"content_type"`
ConversationID string `json:"conversation_id,omitempty"`
ReplyToID string `json:"reply_to_id,omitempty"`
Timestamp int64 `json:"timestamp,omitempty"`
Attachments []FakeAttachment `json:"attachments,omitempty"`
Event string `json:"event"`
MessageID string `json:"message_id"`
SenderID string `json:"sender_id"`
SenderName string `json:"sender_name"`
Content string `json:"content"`
ContentType string `json:"content_type"`
ConversationID string `json:"conversation_id,omitempty"`
ReplyToID string `json:"reply_to_id,omitempty"`
Timestamp int64 `json:"timestamp,omitempty"`
Attachments []FakeAttachment `json:"attachments,omitempty"`
}
// FakeAttachment mirrors channel.Attachment for the fake payload.
@@ -218,11 +218,12 @@ func (p *FakeProvider) ValidateWebhookRequest(ctx context.Context, inbox *model.
// FakeOutboundPayload is the JSON body FakeProvider.SendMessage posts to the
// FakeMessagePlatform /receive endpoint.
type FakeOutboundPayload struct {
MessageID uint `json:"message_id"`
ConversationID uint `json:"conversation_id"`
Content string `json:"content"`
ContentType string `json:"content_type"`
Sender FakeSender `json:"sender"`
MessageID uint `json:"message_id"`
ConversationID uint `json:"conversation_id"`
Content string `json:"content"`
ContentType string `json:"content_type"`
Sender FakeSender `json:"sender"`
Recipient FakeRecipient `json:"recipient"`
}
type FakeSender struct {
@@ -231,6 +232,12 @@ type FakeSender struct {
Type string `json:"type"`
}
type FakeRecipient struct {
ID uint `json:"id"`
Name string `json:"name"`
SourceID string `json:"source_id"`
}
func (p *FakeProvider) SendMessage(ctx context.Context, inbox *model.Inbox, message *model.Message, contact *model.Contact) (*channel.SendResult, error) {
config := parseFakeConfig(inbox)
webhookURL, _ := config["webhook_url"].(string)
@@ -249,9 +256,13 @@ func (p *FakeProvider) SendMessage(ctx context.Context, inbox *model.Inbox, mess
if message.SenderID != nil {
sender.ID = *message.SenderID
}
// Prefer contact name when available, fall back to empty.
recipient := FakeRecipient{}
if contact != nil {
sender.Name = contact.Name
recipient = FakeRecipient{
ID: contact.ID,
Name: contact.Name,
SourceID: contact.Identifier,
}
}
payload := FakeOutboundPayload{
@@ -260,6 +271,7 @@ func (p *FakeProvider) SendMessage(ctx context.Context, inbox *model.Inbox, mess
Content: message.Content,
ContentType: message.ContentType,
Sender: sender,
Recipient: recipient,
}
body, err := json.Marshal(payload)
+24 -10
View File
@@ -207,11 +207,11 @@ func TestFakeProvider_SendMessage_NoWebhookURL(t *testing.T) {
inbox := &model.Inbox{}
inbox.ID = 1
msg := &model.Message{
Base: model.Base{ID: 100},
Base: model.Base{ID: 100},
ConversationID: 50,
Content: "test",
ContentType: "text",
SenderType: "agent",
Content: "test",
ContentType: "text",
SenderType: "agent",
}
result, err := p.SendMessage(context.Background(), inbox, msg, nil)
@@ -247,14 +247,18 @@ func TestFakeProvider_SendMessage_PostsToWebhookURL(t *testing.T) {
senderID := uint(5)
msg := &model.Message{
Base: model.Base{ID: 100},
Base: model.Base{ID: 100},
ConversationID: 50,
Content: "Hello from agent",
ContentType: "text",
SenderType: "agent",
SenderID: &senderID,
Content: "Hello from agent",
ContentType: "text",
SenderType: "agent",
SenderID: &senderID,
}
contact := &model.Contact{
Base: model.Base{ID: 9},
Name: "Customer Chen",
Identifier: "customer_001",
}
contact := &model.Contact{Name: "Agent Wang"}
result, err := p.SendMessage(context.Background(), inbox, msg, contact)
if err != nil {
@@ -278,6 +282,16 @@ func TestFakeProvider_SendMessage_PostsToWebhookURL(t *testing.T) {
if sender["type"] != "agent" {
t.Fatalf("expected sender.type 'agent', got %v", sender["type"])
}
recipient, ok := receivedBody["recipient"].(map[string]interface{})
if !ok {
t.Fatal("expected recipient object")
}
if recipient["source_id"] != "customer_001" {
t.Fatalf("expected recipient.source_id 'customer_001', got %v", recipient["source_id"])
}
if recipient["name"] != "Customer Chen" {
t.Fatalf("expected recipient.name 'Customer Chen', got %v", recipient["name"])
}
}
func TestFakeProvider_GetContactProfile(t *testing.T) {
+2
View File
@@ -760,6 +760,8 @@ func parseDotEnvLine(line string) (string, string, bool) {
val := strings.TrimSpace(parts[1])
if len(val) >= 2 && ((val[0] == '"' && val[len(val)-1] == '"') || (val[0] == '\'' && val[len(val)-1] == '\'')) {
val = val[1 : len(val)-1]
} else if strings.HasPrefix(val, "#") {
val = ""
} else if comment := strings.Index(val, " #"); comment >= 0 {
val = strings.TrimSpace(val[:comment])
}
@@ -4,12 +4,14 @@ import (
"context"
"errors"
"io"
"net"
"net/http"
"strconv"
"strings"
"github.com/gin-gonic/gin"
"github.com/gochat/gochat/internal/llm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/service"
"github.com/gochat/gochat/pkg/pagination"
@@ -1055,6 +1057,29 @@ func handleServiceError(c *gin.Context, err error) {
return
}
errMsg := err.Error()
if errors.Is(err, llm.ErrProviderNotConfigured) {
response.AbortWithStatusError(c, http.StatusServiceUnavailable, response.ErrCopilotNotConfigured, errMsg)
return
}
var providerErr *llm.APIError
if errors.As(err, &providerErr) {
switch providerErr.StatusCode {
case http.StatusUnauthorized, http.StatusForbidden:
response.AbortWithStatusError(c, http.StatusBadGateway, response.ErrCopilotProviderAuth, "Copilot provider authentication failed")
case http.StatusNotFound:
response.AbortWithStatusError(c, http.StatusBadGateway, response.ErrCopilotModelNotFound, "Copilot provider endpoint or model was not found")
case http.StatusTooManyRequests:
response.AbortWithStatusError(c, http.StatusTooManyRequests, response.ErrCopilotProviderRateLimited, "Copilot provider rate limit exceeded")
default:
response.AbortWithStatusError(c, http.StatusBadGateway, response.ErrCopilotProviderUnreachable, "Copilot provider request failed")
}
return
}
var networkErr net.Error
if errors.Is(err, context.DeadlineExceeded) || (errors.As(err, &networkErr) && networkErr.Timeout()) {
response.AbortWithStatusError(c, http.StatusGatewayTimeout, response.ErrCopilotProviderTimeout, "Copilot provider request timed out")
return
}
lower := strings.ToLower(errMsg)
if strings.Contains(lower, "not found") || strings.Contains(lower, "record not found") {
response.AbortWithStatusError(c, http.StatusNotFound, response.ErrNotFound, errMsg)
@@ -694,6 +694,7 @@ func serializeContactWithContext(ctx context.Context, contact *model.Contact) ma
"custom_attributes": jsonObject(contact.CustomAttributes),
"last_activity_at": int64Value(contact.LastActivityAt),
"created_at": contact.CreatedAt.Unix(),
"type": "contact",
}
}
@@ -710,6 +711,7 @@ func serializeUser(user *model.User, accountID uint) map[string]any {
"name": user.Name,
"role": nonEmpty(user.Role, "agent"),
"thumbnail": user.AvatarURL,
"type": "user",
}
if attrs := jsonObject(user.CustomAttributes); len(attrs) > 0 {
payload["custom_attributes"] = attrs
@@ -390,9 +390,19 @@ func TestSerializeContactUsesPresenceStatus(t *testing.T) {
payload := serializeContactWithContext(ctx, contact)
require.Equal(t, "online", payload["availability_status"])
require.Equal(t, "contact", payload["type"])
require.Equal(t, "offline", serializeContact(contact)["availability_status"])
}
func TestSerializeUserIncludesSenderType(t *testing.T) {
user := &model.User{AccountID: 7, Name: "Support Agent"}
user.ID = 43
payload := serializeUser(user, user.AccountID)
require.Equal(t, "user", payload["type"])
}
func TestSerializeCRMContactUsesPresenceStatus(t *testing.T) {
contact := &model.Contact{AccountID: 9, Name: "CRM Online Contact"}
contact.ID = 99
@@ -1759,10 +1759,19 @@ func (s *ConversationService) ToggleTyping(ctx context.Context, accountID, conve
if !ok {
return nil
}
event := channel.NewChannelEvent(eventType, channel.ChannelAPI, accountID, conversation.InboxID)
channelType := channel.ChannelType(conversation.ChannelType)
if channelType == "" {
channelType = channel.ChannelAPI
}
event := channel.NewChannelEvent(eventType, channelType, accountID, conversation.InboxID)
event.ConversationID = conversationID
event.ContactID = conversation.ContactID
event.UserID = userID
event.Data["conversation"] = conversation
var user model.User
if err := s.repo.DB().WithContext(ctx).First(&user, userID).Error; err == nil {
event.Data["user"] = &user
}
event.Data["typing_status"] = typingStatus
event.Data["is_private"] = isPrivate
s.dispatcher.Dispatch(ctx, event)
+188 -5
View File
@@ -9,8 +9,11 @@ package wsevent
import (
"context"
"encoding/json"
"strings"
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
wspkg "github.com/gochat/gochat/internal/ws"
applogger "github.com/gochat/gochat/pkg/logger"
)
@@ -51,13 +54,13 @@ func (l *BridgeListener) OnEvent(ctx context.Context, event *channel.ChannelEven
return nil
}
payload := event.Data
if payload == nil {
payload = map[string]interface{}{}
}
payload := wsEventPayload(event)
payload["account_id"] = event.AccountID
if event.ConversationID != 0 {
payload["conversation_id"] = event.ConversationID
if _, exists := payload["conversation_id"]; !exists {
payload["conversation_id"] = event.ConversationID
}
}
if event.InboxID != 0 {
payload["inbox_id"] = event.InboxID
@@ -69,6 +72,186 @@ func (l *BridgeListener) OnEvent(ctx context.Context, event *channel.ChannelEven
return nil
}
// wsEventPayload converts internal dispatcher data into the flat push payload
// expected by the reused Chatwoot ActionCable client. Chatwoot broadcasts
// message.push_event_data directly, not an internal {message, conversation}
// wrapper. Keeping this normalization at the WS boundary lets other listeners
// continue consuming the richer internal event data.
func wsEventPayload(event *channel.ChannelEvent) map[string]interface{} {
if event == nil {
return map[string]interface{}{}
}
switch event.Type {
case channel.EventMessageCreated, channel.EventMessageUpdated, channel.EventMessageDeleted:
if message, ok := eventMessage(event.Data); ok {
return messagePushPayload(message, event.Data)
}
case channel.EventConversationCreated, channel.EventConversationUpdated,
channel.EventConversationOpened, channel.EventConversationResolved,
channel.EventConversationAssigned, channel.EventConversationUnassigned:
if conversation, ok := eventConversation(event.Data); ok {
return conversationPushPayload(conversation, event.Data)
}
case channel.EventContactCreated, channel.EventContactUpdated, channel.EventContactDeleted:
if contact, ok := eventContact(event.Data); ok {
return contactPushPayload(contact)
}
}
return copyEventData(event.Data)
}
func conversationPushPayload(conversation *model.Conversation, data map[string]interface{}) map[string]interface{} {
payload := modelMap(conversation)
conversationID := conversation.ID
if conversation.DisplayID != nil && *conversation.DisplayID != 0 {
conversationID = *conversation.DisplayID
}
payload["id"] = conversationID
payload["created_at"] = conversation.CreatedAt.Unix()
payload["updated_at"] = float64(conversation.UpdatedAt.UnixNano()) / 1e9
if conversation.LastActivityAt != nil {
payload["last_activity_at"] = *conversation.LastActivityAt
}
if _, exists := payload["messages"]; !exists {
payload["messages"] = []interface{}{}
}
meta := map[string]interface{}{}
if contact, ok := eventContact(data); ok {
meta["sender"] = contactPushPayload(contact)
}
if inbox, ok := eventInbox(data); ok {
meta["channel"] = inbox.ChannelType
}
payload["meta"] = meta
return payload
}
func messagePushPayload(message *model.Message, data map[string]interface{}) map[string]interface{} {
payload := modelMap(message)
payload["created_at"] = message.CreatedAt.Unix()
payload["message_type"] = messageTypeValue(message.MessageType)
payload["content_type"] = nonEmpty(message.ContentType, "text")
payload["status"] = nonEmpty(message.Status, "sent")
if message.ContentAttributes == nil {
payload["content_attributes"] = map[string]interface{}{}
}
conversationID := message.ConversationID
conversationPayload := map[string]interface{}{
"last_activity_at": message.CreatedAt.Unix(),
}
if conversation, ok := eventConversation(data); ok {
if conversation.DisplayID != nil && *conversation.DisplayID != 0 {
conversationID = *conversation.DisplayID
}
conversationPayload["assignee_id"] = conversation.AssigneeID
if conversation.LastActivityAt != nil {
conversationPayload["last_activity_at"] = *conversation.LastActivityAt
}
}
if contact, ok := eventContact(data); ok {
payload["sender"] = contactPushPayload(contact)
conversationPayload["contact_inbox"] = map[string]interface{}{
"source_id": contact.Identifier,
}
}
payload["conversation_id"] = conversationID
payload["conversation"] = conversationPayload
return payload
}
func contactPushPayload(contact *model.Contact) map[string]interface{} {
payload := modelMap(contact)
payload["created_at"] = contact.CreatedAt.Unix()
payload["availability_status"] = "offline"
payload["type"] = "contact"
if contact.AdditionalAttributes == nil {
payload["additional_attributes"] = map[string]interface{}{}
}
if contact.CustomAttributes == nil {
payload["custom_attributes"] = map[string]interface{}{}
}
return payload
}
func eventMessage(data map[string]interface{}) (*model.Message, bool) {
if message, ok := data["message"].(*model.Message); ok && message != nil {
return message, true
}
if message, ok := data["message"].(model.Message); ok {
return &message, true
}
return nil, false
}
func eventConversation(data map[string]interface{}) (*model.Conversation, bool) {
if conversation, ok := data["conversation"].(*model.Conversation); ok && conversation != nil {
return conversation, true
}
if conversation, ok := data["conversation"].(model.Conversation); ok {
return &conversation, true
}
return nil, false
}
func eventContact(data map[string]interface{}) (*model.Contact, bool) {
if contact, ok := data["contact"].(*model.Contact); ok && contact != nil {
return contact, true
}
if contact, ok := data["contact"].(model.Contact); ok {
return &contact, true
}
return nil, false
}
func eventInbox(data map[string]interface{}) (*model.Inbox, bool) {
if inbox, ok := data["inbox"].(*model.Inbox); ok && inbox != nil {
return inbox, true
}
if inbox, ok := data["inbox"].(model.Inbox); ok {
return &inbox, true
}
return nil, false
}
func modelMap(value interface{}) map[string]interface{} {
payload := map[string]interface{}{}
raw, err := json.Marshal(value)
if err != nil {
return payload
}
_ = json.Unmarshal(raw, &payload)
return payload
}
func copyEventData(data map[string]interface{}) map[string]interface{} {
payload := make(map[string]interface{}, len(data))
for key, value := range data {
payload[key] = value
}
return payload
}
func messageTypeValue(value string) int {
switch strings.ToLower(strings.TrimSpace(value)) {
case "incoming":
return 0
case "activity":
return 2
case "template":
return 3
default:
return 1
}
}
func nonEmpty(value, fallback string) string {
if strings.TrimSpace(value) == "" {
return fallback
}
return value
}
// isWSEventType returns true if the event type string corresponds to a
// WebSocket/SSE event constant defined in wspkg.
func isWSEventType(eventType string) bool {
@@ -0,0 +1,148 @@
package wsevent
import (
"context"
"encoding/json"
"testing"
"time"
"github.com/gochat/gochat/internal/channel"
"github.com/gochat/gochat/internal/model"
wspkg "github.com/gochat/gochat/internal/ws"
)
type captureHub struct {
accountID uint
data []byte
}
func (h *captureHub) SendToAccount(accountID uint, data []byte) {
h.accountID = accountID
h.data = append([]byte(nil), data...)
}
func (h *captureHub) SendToRoom(string, []byte) {}
func TestBridgeListenerMessageCreatedUsesChatwootPushPayload(t *testing.T) {
now := time.Unix(1_783_834_149, 0)
lastActivity := now.Unix()
displayID := uint(22)
contactID := uint(9)
message := &model.Message{
Base: model.Base{ID: 12, CreatedAt: now},
AccountID: 1,
InboxID: 4,
ConversationID: 2,
SenderID: &contactID,
SenderType: string(model.SenderTypeContact),
Content: "hellohello",
ContentType: "text",
Status: "sent",
MessageType: "incoming",
SourceID: "fake_echo_11",
}
conversation := &model.Conversation{
Base: model.Base{ID: 2, CreatedAt: now},
AccountID: 1,
InboxID: 4,
ContactID: contactID,
DisplayID: &displayID,
LastActivityAt: &lastActivity,
}
contact := &model.Contact{
Base: model.Base{ID: contactID, CreatedAt: now},
AccountID: 1,
Name: "Fake Customer",
Identifier: "customer_001",
}
event := channel.NewChannelEvent(channel.EventMessageCreated, channel.ChannelFake, 1, 4)
event.ConversationID = conversation.ID
event.Data["message"] = message
event.Data["conversation"] = conversation
event.Data["contact"] = contact
hub := &captureHub{}
listener := New(wspkg.NewEventPublisherLocal(hub, nil))
if err := listener.OnEvent(context.Background(), event); err != nil {
t.Fatalf("OnEvent failed: %v", err)
}
var envelope map[string]interface{}
if err := json.Unmarshal(hub.data, &envelope); err != nil {
t.Fatalf("decode websocket envelope: %v", err)
}
payload, ok := envelope["data"].(map[string]interface{})
if !ok {
t.Fatalf("expected object payload, got %#v", envelope["data"])
}
if payload["account_id"] != float64(1) {
t.Fatalf("expected account_id=1, got %#v", payload["account_id"])
}
if payload["id"] != float64(12) || payload["content"] != "hellohello" {
t.Fatalf("unexpected flattened message payload: %#v", payload)
}
if payload["message_type"] != float64(0) {
t.Fatalf("expected incoming message_type=0, got %#v", payload["message_type"])
}
if payload["conversation_id"] != float64(displayID) {
t.Fatalf("expected display conversation id %d, got %#v", displayID, payload["conversation_id"])
}
if _, nested := payload["message"]; nested {
t.Fatalf("message.created payload must be flat: %#v", payload)
}
sender, ok := payload["sender"].(map[string]interface{})
if !ok || sender["name"] != "Fake Customer" || sender["type"] != "contact" {
t.Fatalf("expected contact sender payload, got %#v", payload["sender"])
}
conversationData, ok := payload["conversation"].(map[string]interface{})
if !ok || conversationData["last_activity_at"] != float64(lastActivity) {
t.Fatalf("expected conversation push payload, got %#v", payload["conversation"])
}
}
func TestBridgeListenerConversationUpdatedIncludesMeta(t *testing.T) {
now := time.Unix(1_783_834_149, 0)
contactID := uint(9)
conversation := &model.Conversation{
Base: model.Base{ID: 2, CreatedAt: now, UpdatedAt: now},
AccountID: 1,
InboxID: 4,
ContactID: contactID,
Status: "open",
}
contact := &model.Contact{
Base: model.Base{ID: contactID, CreatedAt: now},
AccountID: 1,
Name: "Fake Customer",
Identifier: "customer_001",
}
inbox := &model.Inbox{Base: model.Base{ID: 4}, AccountID: 1, ChannelType: "fake"}
event := channel.NewChannelEvent(channel.EventConversationUpdated, channel.ChannelFake, 1, 4)
event.ConversationID = conversation.ID
event.Data["conversation"] = conversation
event.Data["contact"] = contact
event.Data["inbox"] = inbox
hub := &captureHub{}
listener := New(wspkg.NewEventPublisherLocal(hub, nil))
if err := listener.OnEvent(context.Background(), event); err != nil {
t.Fatalf("OnEvent failed: %v", err)
}
var envelope map[string]interface{}
if err := json.Unmarshal(hub.data, &envelope); err != nil {
t.Fatalf("decode websocket envelope: %v", err)
}
payload := envelope["data"].(map[string]interface{})
meta, ok := payload["meta"].(map[string]interface{})
if !ok {
t.Fatalf("expected conversation meta, got %#v", payload["meta"])
}
if meta["channel"] != "fake" {
t.Fatalf("expected fake channel meta, got %#v", meta)
}
sender, ok := meta["sender"].(map[string]interface{})
if !ok || sender["name"] != "Fake Customer" || sender["type"] != "contact" {
t.Fatalf("expected sender in conversation meta, got %#v", meta)
}
}