Align GoChat with Chatwoot frontend contracts

This commit is contained in:
2026-06-13 22:13:32 +08:00
parent 71abf58636
commit d884fdda0a
162 changed files with 9825 additions and 528 deletions
+23 -18
View File
@@ -29,7 +29,7 @@ type ServerMessageType string
const (
// ServerEvent pushes a real-time event to subscribed clients
ServerEvent ServerMessageType = "event"
ServerEvent ServerMessageType = "event"
// ServerConfirmSubscribe acknowledges a successful subscription
ServerConfirmSubscribe ServerMessageType = "confirm_subscribe"
// ServerConfirmUnsubscribe acknowledges a successful unsubscribe
@@ -50,15 +50,15 @@ const (
// Mirrors Chatwoot ActionCable's command structure.
type CommandFrame struct {
Command CommandType `json:"command"`
Identifier string `json:"identifier"` // JSON-encoded ChannelIdentifier
Identifier string `json:"identifier"` // JSON-encoded ChannelIdentifier
Data string `json:"data,omitempty"` // optional action data
}
// ChannelIdentifier describes which "channel" (room) the client wants to subscribe to.
// Serialized as JSON string in the `identifier` field, matching ActionCable convention.
type ChannelIdentifier struct {
Channel string `json:"channel"` // "AccountChannel" or "ConversationChannel"
AccountID uint `json:"account_id"` // required for both channels
Channel string `json:"channel"` // "AccountChannel" or "ConversationChannel"
AccountID uint `json:"account_id"` // required for both channels
ConversationID uint `json:"conversation_id,omitempty"` // required for ConversationChannel
}
@@ -73,9 +73,9 @@ const (
// EventFrame pushes a real-time event payload to the client.
type EventFrame struct {
Type ServerMessageType `json:"type"`
Event string `json:"event,omitempty"` // e.g. "message.created"
Payload interface{} `json:"payload,omitempty"` // event data
Identifier string `json:"identifier,omitempty"` // channel identifier
Event string `json:"event,omitempty"` // e.g. "message.created"
Payload interface{} `json:"payload,omitempty"` // event data
Identifier string `json:"identifier,omitempty"` // channel identifier
}
// ConfirmFrame acknowledges a subscribe/unsubscribe command.
@@ -104,9 +104,9 @@ type WelcomeFrame struct {
// DisconnectFrame is sent before closing a connection.
type DisconnectFrame struct {
Type ServerMessageType `json:"type"`
Reason string `json:"reason"`
Reconnect bool `json:"reconnect"`
Type ServerMessageType `json:"type"`
Reason string `json:"reason"`
Reconnect bool `json:"reconnect"`
}
// --- Real-time Event Type Constants ---
@@ -119,11 +119,11 @@ const (
EventMessageDeleted = "message.deleted"
// Conversation events
EventConversationCreated = "conversation.created"
EventConversationUpdated = "conversation.updated"
EventConversationResolved = "conversation.resolved"
EventConversationOpened = "conversation.opened"
EventConversationAssigned = "conversation.assigned"
EventConversationCreated = "conversation.created"
EventConversationUpdated = "conversation.updated"
EventConversationResolved = "conversation.resolved"
EventConversationOpened = "conversation.opened"
EventConversationAssigned = "conversation.assigned"
EventConversationUnassigned = "conversation.unassigned"
// Contact events
@@ -134,8 +134,8 @@ const (
// Agent/typing events
EventAgentTypingOn = "agent.typing_on"
EventAgentTypingOff = "agent.typing_off"
EventAgentOnline = "agent.online"
EventAgentOffline = "agent.offline"
EventAgentOnline = "agent.online"
EventAgentOffline = "agent.offline"
// Inbox events
EventInboxCreated = "inbox.created"
@@ -147,6 +147,11 @@ const (
// P4 M8 — Notification+Webhook event types
EventNotificationCreated = "notification.created"
EventNotificationUpdated = "notification.updated"
EventNotificationDeleted = "notification.deleted"
// Account cache event types
EventAccountCacheInvalidated = "account.cache_invalidated"
)
// --- Ping/pong Configuration ---
@@ -154,4 +159,4 @@ const (
const (
// PingInterval is how often the server sends ping frames to detect dead connections.
PingInterval = 30 // seconds
)
)
+33 -10
View File
@@ -7,8 +7,8 @@ import (
"strings"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill/message"
"github.com/ThreeDotsLabs/watermill-redisstream/pkg/redisstream"
"github.com/ThreeDotsLabs/watermill/message"
"github.com/redis/go-redis/v9"
"github.com/gochat/gochat/pkg/logger"
@@ -19,17 +19,18 @@ import (
// and forwards the events to the appropriate WebSocket rooms via the Hub.
//
// Architecture mapping:
// Chatwoot ActionCable broadcasts → Watermill subscriber → Hub.SendToAccount/SendToConversationJSON
// Each GoChat instance runs its own subscriber, so events are delivered to
// locally-connected WebSocket clients. Multi-instance delivery relies on Redis
// PubSub (each instance receives the event and pushes to its own clients).
//
// Chatwoot ActionCable broadcasts → Watermill subscriber → Hub.SendToAccount/SendToConversationJSON
// Each GoChat instance runs its own subscriber, so events are delivered to
// locally-connected WebSocket clients. Multi-instance delivery relies on Redis
// PubSub (each instance receives the event and pushes to its own clients).
//
// Reference: P2E §3 — Real-time communication via Redis Pub/Sub
type Subscriber struct {
hub *Hub
subscriber *redisstream.Subscriber
router *message.Router
redisClient redis.UniversalClient
hub *Hub
subscriber *redisstream.Subscriber
router *message.Router
redisClient redis.UniversalClient
}
// NewSubscriber creates a PubSub-to-WebSocket bridge subscriber.
@@ -172,6 +173,28 @@ func (s *Subscriber) registerHandlers() {
s.subscriber,
s.forwardToAccount(EventNotificationCreated),
)
s.router.AddNoPublisherHandler(
"ws-notification-updated-handler",
"gochat.notification.updated",
s.subscriber,
s.forwardToAccount(EventNotificationUpdated),
)
s.router.AddNoPublisherHandler(
"ws-notification-deleted-handler",
"gochat.notification.deleted",
s.subscriber,
s.forwardToAccount(EventNotificationDeleted),
)
// --- Account cache events → Account room ---
s.router.AddNoPublisherHandler(
"ws-account-cache-invalidated-handler",
"gochat.account.cache_invalidated",
s.subscriber,
s.forwardToAccount(EventAccountCacheInvalidated),
)
}
// forwardToAccountAndConversation creates a handler that pushes an event to both
@@ -356,4 +379,4 @@ func (s *Subscriber) Close() error {
// Running returns whether the router is currently running.
func (s *Subscriber) Running() bool {
return s.router.IsRunning()
}
}