Reorganize repo: backend/, deploy/, docs/ layout + AGENTS.md
Restructure the monorepo into clear top-level directories: - backend/: Go module root (cmd, internal, pkg, configs, migrations, docs/swagger, scripts, tests, go.mod, Makefile, .air.toml) - deploy/: Docker (Dockerfile, docker-compose*), quickstart, fluentd - docs/: project documentation + reports/ (moved from repo root) - AGENTS.md: new AI coding-agent guide at repo root Update all references to the new layout: - Dockerfile: COPY backend/go.mod, COPY backend/ (context = repo root) - docker-compose files: context ../.., dockerfile deploy/docker/Dockerfile, env_file ../../.env, volume mounts ../../backend:/app - deploy/quickstart/compose.yaml: dockerfile deploy/docker/Dockerfile - CI: working-directory: backend for go commands, file deploy/docker/Dockerfile, coverage path backend/coverage.out, health_check backend/scripts/ - backend/Makefile: docker target uses -f ../deploy/docker/Dockerfile ../ - README: architecture tree, quickstart, config paths updated Move root stray scripts (rename_models.*, run_m11_tests.sh, verify_build.sh, gorm_bool_main.go) to backend/scripts/legacy/. All moves via git mv to preserve history. Build, vet, SQLite tests, and docker compose config verified.
This commit is contained in:
@@ -0,0 +1,289 @@
|
||||
package ws
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/gorilla/websocket"
|
||||
|
||||
wspkg "github.com/gochat/gochat/internal/ws"
|
||||
"github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
// uintToStr converts a uint to its string representation.
|
||||
// Kept for potential use in identifier formatting.
|
||||
func uintToStr(u uint) string {
|
||||
return strconv.FormatUint(uint64(u), 10)
|
||||
}
|
||||
|
||||
// Handler manages WebSocket connections, upgrade, JWT auth, and message dispatch.
|
||||
// Reference: Chatwoot ActionCable — WebSocket upgrade at /cable, JWT-based auth,
|
||||
// room-based subscription model (AccountChannel, ConversationChannel).
|
||||
//
|
||||
// Supports two authentication paths:
|
||||
// 1. Agent/User auth: JWT token (from query param or Authorization header)
|
||||
// 2. Contact auth: pubsub_token + user_id (Chatwoot RoomChannel pattern)
|
||||
type Handler struct {
|
||||
hub *Hub
|
||||
authenticator *wspkg.WSAuthenticator
|
||||
upgrader websocket.Upgrader
|
||||
}
|
||||
|
||||
// NewHandler creates a WebSocket handler with the given hub and authenticator.
|
||||
// The authenticator provides both JWT (agent) and pubsub_token (contact) auth paths.
|
||||
func NewHandler(hub *Hub, authenticator *wspkg.WSAuthenticator) *Handler {
|
||||
return &Handler{
|
||||
hub: hub,
|
||||
authenticator: authenticator,
|
||||
upgrader: websocket.Upgrader{
|
||||
ReadBufferSize: 1024,
|
||||
WriteBufferSize: 1024,
|
||||
// Allow all origins — CORS is handled at the Gin middleware layer
|
||||
CheckOrigin: func(r *http.Request) bool { return true },
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// ServeWS handles the WebSocket upgrade request at /ws.
|
||||
// URL: /ws?token=<JWT_ACCESS_TOKEN> OR /ws?pubsub_token=<TOKEN>&user_id=<ID>
|
||||
// On successful upgrade, the handler:
|
||||
// 1. Authenticates the request via WSAuthenticator (JWT or pubsub_token)
|
||||
// 2. Authorizes the user/contact for the requested account
|
||||
// 3. Upgrades HTTP connection to WebSocket
|
||||
// 4. Creates a Client and registers it with the Hub
|
||||
// 5. Starts readPump (incoming commands) and writePump (outgoing events) goroutines
|
||||
func (h *Handler) ServeWS(c *gin.Context) {
|
||||
// Step 1: Authenticate (JWT or pubsub_token)
|
||||
claims, err := h.authenticator.Authenticate(c)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: authentication failed: %v", err)
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Step 2: Authorize (verify account access)
|
||||
if err := h.authenticator.Authorize(claims, c); err != nil {
|
||||
logger.L().Errorf("ws: authorization failed: %v", err)
|
||||
c.JSON(http.StatusForbidden, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Step 3: Upgrade HTTP connection to WebSocket
|
||||
conn, err := h.upgrader.Upgrade(c.Writer, c.Request, nil)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: upgrade failed: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Step 4: Create client and register with hub
|
||||
client := NewClient(claims.UserID, claims.AccountID, conn, h.hub)
|
||||
// Populate extended claims fields on the client
|
||||
client.Role = claims.Role
|
||||
client.IsContact = claims.IsContact
|
||||
client.PubsubToken = claims.PubsubToken
|
||||
client.ContactID = claims.ContactID
|
||||
client.InboxID = claims.InboxID
|
||||
|
||||
h.hub.Register(client)
|
||||
|
||||
logger.L().Infof("ws: connection established (user=%d, account=%d, is_contact=%v)", claims.UserID, claims.AccountID, claims.IsContact)
|
||||
|
||||
// Start pumps in separate goroutines
|
||||
go h.writePump(client)
|
||||
go h.readPump(client)
|
||||
}
|
||||
|
||||
// ServeCable handles the WebSocket upgrade request at /cable.
|
||||
// This is the ActionCable-compatible endpoint (Chatwoot uses /cable for WS).
|
||||
// It is functionally identical to ServeWS but uses the ActionCable naming convention.
|
||||
// Some Chatwoot clients specifically connect to /cable.
|
||||
func (h *Handler) ServeCable(c *gin.Context) {
|
||||
h.ServeWS(c)
|
||||
}
|
||||
|
||||
// readPump reads messages from the WebSocket connection and dispatches commands.
|
||||
// Reference: Chatwoot ActionCable consumer — processes subscribe/unsubscribe/ping commands.
|
||||
// One readPump runs per client connection.
|
||||
func (h *Handler) readPump(client *Client) {
|
||||
defer func() {
|
||||
h.hub.Unregister(client)
|
||||
client.Conn.Close()
|
||||
}()
|
||||
|
||||
client.Conn.SetReadLimit(MaxMessageSize)
|
||||
client.Conn.SetReadDeadline(time.Now().Add(PongWait))
|
||||
client.Conn.SetPongHandler(func(string) error {
|
||||
client.Conn.SetReadDeadline(time.Now().Add(PongWait))
|
||||
return nil
|
||||
})
|
||||
|
||||
for {
|
||||
_, message, err := client.Conn.ReadMessage()
|
||||
if err != nil {
|
||||
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
|
||||
logger.L().Errorf("ws: unexpected close for user=%d: %v", client.UserID, err)
|
||||
}
|
||||
break
|
||||
}
|
||||
|
||||
// Decode the command frame
|
||||
var cmd CommandFrame
|
||||
if err := json.Unmarshal(message, &cmd); err != nil {
|
||||
logger.L().Warnf("ws: invalid command from user=%d: %v", client.UserID, err)
|
||||
continue
|
||||
}
|
||||
|
||||
// Dispatch command
|
||||
switch cmd.Command {
|
||||
case CommandSubscribe:
|
||||
h.handleSubscribe(client, cmd)
|
||||
case CommandUnsubscribe:
|
||||
h.handleUnsubscribe(client, cmd)
|
||||
case CommandPing:
|
||||
h.handlePing(client)
|
||||
default:
|
||||
logger.L().Warnf("ws: unknown command '%s' from user=%d", cmd.Command, client.UserID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// writePump sends messages from the client's Send channel to the WebSocket connection.
|
||||
// It also sends periodic ping frames for heartbeat detection.
|
||||
// One writePump runs per client connection.
|
||||
func (h *Handler) writePump(client *Client) {
|
||||
ticker := time.NewTicker(time.Duration(PingInterval) * time.Second)
|
||||
defer func() {
|
||||
ticker.Stop()
|
||||
client.Conn.Close()
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case message, ok := <-client.Send:
|
||||
client.Conn.SetWriteDeadline(time.Now().Add(WriteWait))
|
||||
if !ok {
|
||||
// Hub closed the channel — send close frame
|
||||
client.Conn.WriteMessage(websocket.CloseMessage, []byte{})
|
||||
return
|
||||
}
|
||||
|
||||
// Write text message (all our frames are JSON text)
|
||||
if err := client.Conn.WriteMessage(websocket.TextMessage, message); err != nil {
|
||||
logger.L().Errorf("ws: write error for user=%d: %v", client.UserID, err)
|
||||
return
|
||||
}
|
||||
|
||||
case <-ticker.C:
|
||||
// Send ping frame for heartbeat
|
||||
client.Conn.SetWriteDeadline(time.Now().Add(WriteWait))
|
||||
if err := client.Conn.WriteMessage(websocket.PingMessage, []byte{}); err != nil {
|
||||
logger.L().Errorf("ws: ping write failed for user=%d: %v", client.UserID, err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// handleSubscribe processes a subscribe command.
|
||||
// Validates the ChannelIdentifier and adds the client to the appropriate room.
|
||||
func (h *Handler) handleSubscribe(client *Client, cmd CommandFrame) {
|
||||
var identifier ChannelIdentifier
|
||||
if err := json.Unmarshal([]byte(cmd.Identifier), &identifier); err != nil {
|
||||
logger.L().Warnf("ws: invalid identifier from user=%d: %v", client.UserID, err)
|
||||
// Send reject frame
|
||||
rejectData, _ := json.Marshal(RejectFrame{
|
||||
Type: ServerRejectSubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
Reason: "invalid identifier format",
|
||||
})
|
||||
client.Send <- rejectData
|
||||
return
|
||||
}
|
||||
|
||||
// Validate account_id matches the client's authenticated account
|
||||
if identifier.AccountID != client.AccountID {
|
||||
rejectData, _ := json.Marshal(RejectFrame{
|
||||
Type: ServerRejectSubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
Reason: "account_id mismatch",
|
||||
})
|
||||
client.Send <- rejectData
|
||||
return
|
||||
}
|
||||
|
||||
// Determine room name based on channel type (uses Hub's canonical naming)
|
||||
room := ""
|
||||
switch identifier.Channel {
|
||||
case ChannelAccount:
|
||||
room = accountRoomName(identifier.AccountID)
|
||||
case ChannelConversation:
|
||||
if identifier.ConversationID == 0 {
|
||||
rejectData, _ := json.Marshal(RejectFrame{
|
||||
Type: ServerRejectSubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
Reason: "conversation_id required for ConversationChannel",
|
||||
})
|
||||
client.Send <- rejectData
|
||||
return
|
||||
}
|
||||
room = conversationRoomName(identifier.AccountID, identifier.ConversationID)
|
||||
default:
|
||||
rejectData, _ := json.Marshal(RejectFrame{
|
||||
Type: ServerRejectSubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
Reason: "unknown channel type: " + identifier.Channel,
|
||||
})
|
||||
client.Send <- rejectData
|
||||
return
|
||||
}
|
||||
|
||||
// Subscribe the client to the room
|
||||
client.Subscribe(room)
|
||||
|
||||
// Send confirmation frame
|
||||
confirmData, _ := json.Marshal(ConfirmFrame{
|
||||
Type: ServerConfirmSubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
})
|
||||
client.Send <- confirmData
|
||||
}
|
||||
|
||||
// handleUnsubscribe processes an unsubscribe command.
|
||||
func (h *Handler) handleUnsubscribe(client *Client, cmd CommandFrame) {
|
||||
var identifier ChannelIdentifier
|
||||
if err := json.Unmarshal([]byte(cmd.Identifier), &identifier); err != nil {
|
||||
logger.L().Warnf("ws: invalid identifier in unsubscribe from user=%d: %v", client.UserID, err)
|
||||
return
|
||||
}
|
||||
|
||||
// Determine room name (uses Hub's canonical naming)
|
||||
room := ""
|
||||
switch identifier.Channel {
|
||||
case ChannelAccount:
|
||||
room = accountRoomName(identifier.AccountID)
|
||||
case ChannelConversation:
|
||||
room = conversationRoomName(identifier.AccountID, identifier.ConversationID)
|
||||
default:
|
||||
return
|
||||
}
|
||||
|
||||
client.Unsubscribe(room)
|
||||
|
||||
confirmData, _ := json.Marshal(ConfirmFrame{
|
||||
Type: ServerConfirmUnsubscribe,
|
||||
Identifier: cmd.Identifier,
|
||||
})
|
||||
client.Send <- confirmData
|
||||
}
|
||||
|
||||
// handlePing responds to a client-initiated ping command with a pong frame.
|
||||
func (h *Handler) handlePing(client *Client) {
|
||||
pingData, _ := json.Marshal(PingFrame{
|
||||
Type: ServerPing,
|
||||
Message: time.Now().UTC().Format(time.RFC3339),
|
||||
})
|
||||
client.Send <- pingData
|
||||
}
|
||||
@@ -0,0 +1,582 @@
|
||||
// Package ws (handler/ws) provides the WebSocket Hub that manages connected
|
||||
// clients, room subscriptions, and message routing. Enhanced with Redis
|
||||
// Pub/Sub relay for cross-instance broadcasting.
|
||||
package ws
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
|
||||
wspkg "github.com/gochat/gochat/internal/ws"
|
||||
"github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
// Client represents a connected WebSocket client.
|
||||
// Each client has a unique ID, user identity from auth claims,
|
||||
// and a buffered Send channel for outgoing messages.
|
||||
type Client struct {
|
||||
ID string // unique connection ID (uuid)
|
||||
UserID uint // from WSClaims.UserID
|
||||
AccountID uint // from WSClaims.AccountID
|
||||
Role string // from WSClaims.Role
|
||||
IsContact bool // from WSClaims.IsContact
|
||||
PubsubToken string // from WSClaims.PubsubToken
|
||||
ContactID uint // from WSClaims.ContactID (only for contacts)
|
||||
InboxID uint // from WSClaims.InboxID (only for contacts)
|
||||
Conn *websocket.Conn // gorilla/websocket connection
|
||||
Send chan []byte // buffered outgoing message channel (256 capacity)
|
||||
Hub *Hub // reference back to Hub
|
||||
SubscribedRooms map[string]bool // rooms this client is subscribed to
|
||||
CancelPresence context.CancelFunc // cancel presence refresh on disconnect
|
||||
}
|
||||
|
||||
// NewClient creates a new WebSocket client with the given identity and connection.
|
||||
func NewClient(userID, accountID uint, conn *websocket.Conn, hub *Hub) *Client {
|
||||
return &Client{
|
||||
ID: generateClientID(),
|
||||
UserID: userID,
|
||||
AccountID: accountID,
|
||||
Conn: conn,
|
||||
Send: make(chan []byte, SendChannelSize),
|
||||
Hub: hub,
|
||||
SubscribedRooms: make(map[string]bool),
|
||||
}
|
||||
}
|
||||
|
||||
// Subscribe adds the client to a room and registers with the Hub.
|
||||
func (c *Client) Subscribe(room string) {
|
||||
c.Hub.mu.Lock()
|
||||
c.SubscribedRooms[room] = true
|
||||
c.Hub.subscribeClient(c.ID, room)
|
||||
c.Hub.mu.Unlock()
|
||||
}
|
||||
|
||||
// Unsubscribe removes the client from a room and unregisters with the Hub.
|
||||
func (c *Client) Unsubscribe(room string) {
|
||||
c.Hub.mu.Lock()
|
||||
delete(c.SubscribedRooms, room)
|
||||
c.Hub.unsubscribeClient(c.ID, room)
|
||||
c.Hub.mu.Unlock()
|
||||
}
|
||||
|
||||
// generateClientID creates a unique client ID using crypto/rand.
|
||||
func generateClientID() string {
|
||||
b := make([]byte, 16)
|
||||
_, _ = rand.Read(b)
|
||||
return hex.EncodeToString(b)
|
||||
}
|
||||
|
||||
const (
|
||||
// MaxMessageSize is the maximum size of a message that can be read from a client.
|
||||
MaxMessageSize = 512 * 1024 // 512KB
|
||||
|
||||
// SendChannelSize is the buffer size for client Send channels.
|
||||
SendChannelSize = 256
|
||||
|
||||
// WriteWait is the time allowed to write a message to the peer.
|
||||
WriteWait = 10 * time.Second
|
||||
|
||||
// PongWait is the time allowed to read the next pong message from the peer.
|
||||
PongWait = 60 * time.Second
|
||||
|
||||
// PingPeriod is how often pings are sent to peers. Must be < PongWait.
|
||||
PingPeriod = (PongWait * 9) / 10
|
||||
)
|
||||
|
||||
// Hub maintains the set of active clients and broadcasts messages to them.
|
||||
// Enhanced with Redis Pub/Sub relay for cross-instance message delivery,
|
||||
// typing indicator support, and agent/contact presence tracking.
|
||||
type Hub struct {
|
||||
// Registered clients indexed by connection ID
|
||||
clients map[string]*Client
|
||||
|
||||
// Room subscriptions: room name → set of client IDs in that room
|
||||
rooms map[string]map[string]bool
|
||||
|
||||
// Account subscriptions: account ID → set of client IDs
|
||||
accounts map[uint]map[string]bool
|
||||
|
||||
// Broadcast relay for cross-instance message delivery via Redis Pub/Sub
|
||||
relay *wspkg.BroadcastRelay
|
||||
|
||||
// Typing tracker for typing indicator management
|
||||
typing *wspkg.TypingTracker
|
||||
|
||||
// Presence tracker for online/offline status
|
||||
presence *wspkg.PresenceTracker
|
||||
|
||||
// Presence lifecycle manager
|
||||
presenceMgr *wspkg.PresenceManager
|
||||
|
||||
// Inbound messages from clients (commands)
|
||||
commandChan chan *ClientCommand
|
||||
|
||||
mu sync.RWMutex
|
||||
}
|
||||
|
||||
// ClientCommand wraps a command from a client with the client reference.
|
||||
type ClientCommand struct {
|
||||
Client *Client
|
||||
Cmd wspkg.WSCommand
|
||||
}
|
||||
|
||||
// NewHub creates a new Hub with Redis Pub/Sub relay and presence subsystems.
|
||||
func NewHub(relay *wspkg.BroadcastRelay, typing *wspkg.TypingTracker,
|
||||
presence *wspkg.PresenceTracker, presenceMgr *wspkg.PresenceManager) *Hub {
|
||||
return &Hub{
|
||||
clients: make(map[string]*Client),
|
||||
rooms: make(map[string]map[string]bool),
|
||||
accounts: make(map[uint]map[string]bool),
|
||||
relay: relay,
|
||||
typing: typing,
|
||||
presence: presence,
|
||||
presenceMgr: presenceMgr,
|
||||
commandChan: make(chan *ClientCommand, 256),
|
||||
}
|
||||
}
|
||||
|
||||
// NewHubSimple creates a minimal Hub without Redis subsystems.
|
||||
// Used for development/testing when Redis is not available.
|
||||
func NewHubSimple() *Hub {
|
||||
return &Hub{
|
||||
clients: make(map[string]*Client),
|
||||
rooms: make(map[string]map[string]bool),
|
||||
accounts: make(map[uint]map[string]bool),
|
||||
commandChan: make(chan *ClientCommand, 256),
|
||||
}
|
||||
}
|
||||
|
||||
// Run starts the Hub's main event loop. Processes client registrations,
|
||||
// unregistrations, and commands. Must be called in a goroutine.
|
||||
func (h *Hub) Run(ctx context.Context) {
|
||||
logger.L().Info("ws hub: starting event loop")
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
logger.L().Info("ws hub: shutting down event loop")
|
||||
h.shutdown()
|
||||
return
|
||||
case cmd := <-h.commandChan:
|
||||
h.processCommand(cmd)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Register adds a client to the Hub and sets up presence.
|
||||
func (h *Hub) Register(c *Client) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
h.clients[c.ID] = c
|
||||
|
||||
// Auto-subscribe to account room on connect
|
||||
roomName := accountRoomName(c.AccountID)
|
||||
h.subscribeClient(c.ID, roomName)
|
||||
c.SubscribedRooms[roomName] = true
|
||||
|
||||
// Set up presence tracking
|
||||
if h.presenceMgr != nil {
|
||||
if c.IsContact {
|
||||
h.presenceMgr.OnContactConnect(context.Background(), c.ContactID, c.AccountID)
|
||||
// Also auto-subscribe to pubsub_token room (Chatwoot RoomChannel pattern)
|
||||
if c.PubsubToken != "" {
|
||||
tokenRoom := pubsubTokenRoomName(c.PubsubToken)
|
||||
h.subscribeClient(c.ID, tokenRoom)
|
||||
c.SubscribedRooms[tokenRoom] = true
|
||||
}
|
||||
} else {
|
||||
cancelFn := h.presenceMgr.OnAgentConnect(context.Background(), c.UserID, c.AccountID)
|
||||
c.CancelPresence = cancelFn
|
||||
}
|
||||
}
|
||||
|
||||
logger.L().Infof("ws hub: client registered (id=%s, user_id=%d, account_id=%d, is_contact=%v)",
|
||||
c.ID, c.UserID, c.AccountID, c.IsContact)
|
||||
|
||||
// Send welcome message
|
||||
welcomeMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventWelcome,
|
||||
Data: map[string]any{"client_id": c.ID},
|
||||
})
|
||||
c.Send <- welcomeMsg
|
||||
}
|
||||
|
||||
// Unregister removes a client from the Hub and cleans up presence.
|
||||
func (h *Hub) Unregister(c *Client) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if _, ok := h.clients[c.ID]; !ok {
|
||||
return // already unregistered
|
||||
}
|
||||
|
||||
// Remove from all subscribed rooms
|
||||
for room := range c.SubscribedRooms {
|
||||
h.unsubscribeClient(c.ID, room)
|
||||
}
|
||||
|
||||
// Remove from Hub's client map
|
||||
delete(h.clients, c.ID)
|
||||
|
||||
// Clean up presence tracking
|
||||
if h.presenceMgr != nil {
|
||||
if c.IsContact {
|
||||
h.presenceMgr.OnContactDisconnect(context.Background(), c.ContactID, c.AccountID)
|
||||
} else {
|
||||
if c.CancelPresence != nil {
|
||||
c.CancelPresence() // stop presence refresh loop
|
||||
}
|
||||
h.presenceMgr.OnAgentDisconnect(context.Background(), c.UserID, c.AccountID)
|
||||
}
|
||||
}
|
||||
|
||||
// Close Send channel
|
||||
close(c.Send)
|
||||
|
||||
logger.L().Infof("ws hub: client unregistered (id=%s, user_id=%d)", c.ID, c.UserID)
|
||||
}
|
||||
|
||||
// SendToAccount sends a message to all clients subscribed to an account room.
|
||||
func (h *Hub) SendToAccount(accountID uint, data []byte) {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
roomName := accountRoomName(accountID)
|
||||
if clientIDs, ok := h.rooms[roomName]; ok {
|
||||
for clientID := range clientIDs {
|
||||
if client, ok := h.clients[clientID]; ok {
|
||||
select {
|
||||
case client.Send <- data:
|
||||
default:
|
||||
logger.L().Warnf("ws hub: dropping message for slow client %s", clientID)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SendToAccountConversation sends a message to all clients subscribed to
|
||||
// a specific conversation within an account.
|
||||
func (h *Hub) SendToAccountConversation(accountID uint, conversationID uint, data []byte) {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
roomName := conversationRoomName(accountID, conversationID)
|
||||
if clientIDs, ok := h.rooms[roomName]; ok {
|
||||
for clientID := range clientIDs {
|
||||
if client, ok := h.clients[clientID]; ok {
|
||||
select {
|
||||
case client.Send <- data:
|
||||
default:
|
||||
logger.L().Warnf("ws hub: dropping message for slow client %s", clientID)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SendToRoom sends a message to all clients in a named room.
|
||||
func (h *Hub) SendToRoom(room string, data []byte) {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
if clientIDs, ok := h.rooms[room]; ok {
|
||||
for clientID := range clientIDs {
|
||||
if client, ok := h.clients[clientID]; ok {
|
||||
select {
|
||||
case client.Send <- data:
|
||||
default:
|
||||
logger.L().Warnf("ws hub: dropping message for slow client %s", clientID)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SendToClient sends a message to a specific client by ID.
|
||||
func (h *Hub) SendToClient(clientID string, data []byte) {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
if client, ok := h.clients[clientID]; ok {
|
||||
select {
|
||||
case client.Send <- data:
|
||||
default:
|
||||
logger.L().Warnf("ws hub: dropping message for slow client %s", clientID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SubmitCommand enqueues a client command for the Hub's event loop to process.
|
||||
func (h *Hub) SubmitCommand(client *Client, cmd wspkg.WSCommand) {
|
||||
h.commandChan <- &ClientCommand{Client: client, Cmd: cmd}
|
||||
}
|
||||
|
||||
// processCommand handles a client command in the Hub's event loop.
|
||||
func (h *Hub) processCommand(cmd *ClientCommand) {
|
||||
switch cmd.Cmd.Command {
|
||||
case "subscribe":
|
||||
h.handleSubscribe(cmd)
|
||||
case "unsubscribe":
|
||||
h.handleUnsubscribe(cmd)
|
||||
case "ping":
|
||||
h.handlePing(cmd)
|
||||
case "typing_on":
|
||||
h.handleTypingOn(cmd)
|
||||
case "typing_off":
|
||||
h.handleTypingOff(cmd)
|
||||
case "update_presence":
|
||||
h.handleUpdatePresence(cmd)
|
||||
default:
|
||||
logger.L().Warnf("ws hub: unknown command '%s' from client %s", cmd.Cmd.Command, cmd.Client.ID)
|
||||
}
|
||||
}
|
||||
|
||||
// handleSubscribe processes a subscribe command.
|
||||
// Mirrors Chatwoot's RoomChannel.subscribe flow.
|
||||
func (h *Hub) handleSubscribe(cmd *ClientCommand) {
|
||||
var data wspkg.SubscribeData
|
||||
if err := json.Unmarshal([]byte(cmd.Cmd.Data), &data); err != nil {
|
||||
logger.L().Warnf("ws hub: invalid subscribe data from client %s: %v", cmd.Client.ID, err)
|
||||
return
|
||||
}
|
||||
|
||||
// Verify account access
|
||||
if data.AccountID != cmd.Client.AccountID {
|
||||
rejectMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventSubscribeReject,
|
||||
Data: map[string]any{"reason": "account_id mismatch"},
|
||||
})
|
||||
cmd.Client.Send <- rejectMsg
|
||||
return
|
||||
}
|
||||
|
||||
var roomName string
|
||||
switch data.Channel {
|
||||
case wspkg.ChannelAccount:
|
||||
roomName = accountRoomName(data.AccountID)
|
||||
case wspkg.ChannelConversation:
|
||||
if data.ConversationID == 0 {
|
||||
rejectMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventSubscribeReject,
|
||||
Data: map[string]any{"reason": "conversation_id required for ConversationChannel"},
|
||||
})
|
||||
cmd.Client.Send <- rejectMsg
|
||||
return
|
||||
}
|
||||
roomName = conversationRoomName(data.AccountID, data.ConversationID)
|
||||
default:
|
||||
rejectMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventSubscribeReject,
|
||||
Data: map[string]any{"reason": "unknown channel type"},
|
||||
})
|
||||
cmd.Client.Send <- rejectMsg
|
||||
return
|
||||
}
|
||||
|
||||
h.mu.Lock()
|
||||
h.subscribeClient(cmd.Client.ID, roomName)
|
||||
cmd.Client.SubscribedRooms[roomName] = true
|
||||
h.mu.Unlock()
|
||||
|
||||
confirmMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventSubscribeConfirm,
|
||||
Data: map[string]any{"room": roomName, "channel": data.Channel},
|
||||
})
|
||||
cmd.Client.Send <- confirmMsg
|
||||
|
||||
logger.L().Infof("ws hub: client %s subscribed to room %s", cmd.Client.ID, roomName)
|
||||
}
|
||||
|
||||
// handleUnsubscribe processes an unsubscribe command.
|
||||
func (h *Hub) handleUnsubscribe(cmd *ClientCommand) {
|
||||
var data wspkg.SubscribeData
|
||||
if err := json.Unmarshal([]byte(cmd.Cmd.Data), &data); err != nil {
|
||||
logger.L().Warnf("ws hub: invalid unsubscribe data from client %s: %v", cmd.Client.ID, err)
|
||||
return
|
||||
}
|
||||
|
||||
var roomName string
|
||||
switch data.Channel {
|
||||
case wspkg.ChannelAccount:
|
||||
roomName = accountRoomName(data.AccountID)
|
||||
case wspkg.ChannelConversation:
|
||||
roomName = conversationRoomName(data.AccountID, data.ConversationID)
|
||||
default:
|
||||
return
|
||||
}
|
||||
|
||||
h.mu.Lock()
|
||||
h.unsubscribeClient(cmd.Client.ID, roomName)
|
||||
delete(cmd.Client.SubscribedRooms, roomName)
|
||||
h.mu.Unlock()
|
||||
|
||||
confirmMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventUnsubscribeConfirm,
|
||||
Data: map[string]any{"room": roomName},
|
||||
})
|
||||
cmd.Client.Send <- confirmMsg
|
||||
|
||||
logger.L().Infof("ws hub: client %s unsubscribed from room %s", cmd.Client.ID, roomName)
|
||||
}
|
||||
|
||||
// handlePing responds to a ping command from the client.
|
||||
func (h *Hub) handlePing(cmd *ClientCommand) {
|
||||
pingMsg, _ := json.Marshal(wspkg.WSMessage{
|
||||
Event: wspkg.EventPingResponse,
|
||||
Data: map[string]any{"timestamp": time.Now().UnixMilli()},
|
||||
})
|
||||
cmd.Client.Send <- pingMsg
|
||||
}
|
||||
|
||||
// handleTypingOn processes a typing_on command.
|
||||
func (h *Hub) handleTypingOn(cmd *ClientCommand) {
|
||||
if h.typing == nil {
|
||||
return
|
||||
}
|
||||
|
||||
var data wspkg.TypingData
|
||||
if err := json.Unmarshal([]byte(cmd.Cmd.Data), &data); err != nil {
|
||||
logger.L().Warnf("ws hub: invalid typing_on data from client %s: %v", cmd.Client.ID, err)
|
||||
return
|
||||
}
|
||||
|
||||
performerType := "user"
|
||||
if cmd.Client.IsContact {
|
||||
performerType = "contact"
|
||||
}
|
||||
|
||||
performer := &wspkg.Performer{
|
||||
ID: cmd.Client.UserID,
|
||||
Type: performerType,
|
||||
}
|
||||
|
||||
if err := h.typing.SetTypingOn(context.Background(), data.AccountID, data.ConversationID, performer); err != nil {
|
||||
logger.L().Warnf("ws hub: typing_on failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// handleTypingOff processes a typing_off command.
|
||||
func (h *Hub) handleTypingOff(cmd *ClientCommand) {
|
||||
if h.typing == nil {
|
||||
return
|
||||
}
|
||||
|
||||
var data wspkg.TypingData
|
||||
if err := json.Unmarshal([]byte(cmd.Cmd.Data), &data); err != nil {
|
||||
logger.L().Warnf("ws hub: invalid typing_off data from client %s: %v", cmd.Client.ID, err)
|
||||
return
|
||||
}
|
||||
|
||||
performerType := "user"
|
||||
if cmd.Client.IsContact {
|
||||
performerType = "contact"
|
||||
}
|
||||
|
||||
performer := &wspkg.Performer{
|
||||
ID: cmd.Client.UserID,
|
||||
Type: performerType,
|
||||
}
|
||||
|
||||
if err := h.typing.SetTypingOff(context.Background(), data.AccountID, data.ConversationID, performer); err != nil {
|
||||
logger.L().Warnf("ws hub: typing_off failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// handleUpdatePresence processes an update_presence command.
|
||||
func (h *Hub) handleUpdatePresence(cmd *ClientCommand) {
|
||||
if h.presence == nil {
|
||||
return
|
||||
}
|
||||
|
||||
var data wspkg.PresenceData
|
||||
if err := json.Unmarshal([]byte(cmd.Cmd.Data), &data); err != nil {
|
||||
logger.L().Warnf("ws hub: invalid update_presence data from client %s: %v", cmd.Client.ID, err)
|
||||
return
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
switch data.Status {
|
||||
case "online":
|
||||
h.presence.SetAgentOnline(ctx, cmd.Client.UserID, cmd.Client.AccountID)
|
||||
case "busy":
|
||||
h.presence.SetAgentBusy(ctx, cmd.Client.UserID, cmd.Client.AccountID)
|
||||
case "offline":
|
||||
h.presence.SetAgentOffline(ctx, cmd.Client.UserID, cmd.Client.AccountID)
|
||||
default:
|
||||
logger.L().Warnf("ws hub: unknown presence status '%s' from client %s", data.Status, cmd.Client.ID)
|
||||
}
|
||||
}
|
||||
|
||||
// --- Internal helpers ---
|
||||
|
||||
// subscribeClient adds a client ID to a room's member set.
|
||||
// Must be called with h.mu held.
|
||||
func (h *Hub) subscribeClient(clientID string, room string) {
|
||||
if _, ok := h.rooms[room]; !ok {
|
||||
h.rooms[room] = make(map[string]bool)
|
||||
}
|
||||
h.rooms[room][clientID] = true
|
||||
}
|
||||
|
||||
// unsubscribeClient removes a client ID from a room's member set.
|
||||
// Must be called with h.mu held.
|
||||
func (h *Hub) unsubscribeClient(clientID string, room string) {
|
||||
if clients, ok := h.rooms[room]; ok {
|
||||
delete(clients, clientID)
|
||||
if len(clients) == 0 {
|
||||
delete(h.rooms, room)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// shutdown disconnects all clients and cleans up.
|
||||
func (h *Hub) shutdown() {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
for id, client := range h.clients {
|
||||
if client.CancelPresence != nil {
|
||||
client.CancelPresence()
|
||||
}
|
||||
close(client.Send)
|
||||
client.Conn.Close()
|
||||
delete(h.clients, id)
|
||||
}
|
||||
|
||||
h.rooms = make(map[string]map[string]bool)
|
||||
h.accounts = make(map[uint]map[string]bool)
|
||||
|
||||
logger.L().Info("ws hub: all clients disconnected, shutdown complete")
|
||||
}
|
||||
|
||||
// Shutdown gracefully closes the Hub, disconnecting all WebSocket clients.
|
||||
// This is the public API called during application graceful shutdown.
|
||||
// It sends a close message to each client and waits for them to finish,
|
||||
// respecting the context deadline.
|
||||
func (h *Hub) Shutdown(ctx context.Context) {
|
||||
logger.L().Info("ws hub: initiating graceful shutdown")
|
||||
h.shutdown()
|
||||
}
|
||||
|
||||
// Room name helpers following Chatwoot naming conventions.
|
||||
|
||||
func accountRoomName(accountID uint) string {
|
||||
return fmt.Sprintf("account_%d", accountID)
|
||||
}
|
||||
|
||||
func conversationRoomName(accountID uint, conversationID uint) string {
|
||||
return fmt.Sprintf("account_%d_conversation_%d", accountID, conversationID)
|
||||
}
|
||||
|
||||
func pubsubTokenRoomName(token string) string {
|
||||
return fmt.Sprintf("pubsub_token_%s", token)
|
||||
}
|
||||
@@ -0,0 +1,162 @@
|
||||
package ws
|
||||
|
||||
// Protocol defines the WebSocket message frame format and event type constants.
|
||||
// Reference: Chatwoot ActionCable protocol — subscribe/command/message pattern.
|
||||
//
|
||||
// Client → Server commands:
|
||||
// {"command":"subscribe","identifier":"{\"channel\":\"AccountChannel\",\"account_id\":1}"}}
|
||||
// {"command":"unsubscribe","identifier":"{\"channel\":\"AccountChannel\",\"account_id\":1}"}}
|
||||
// {"command":"ping"}
|
||||
//
|
||||
// Server → Client events:
|
||||
// {"type":"event","event":"message.created","payload":{...},"identifier":"{\"channel\":\"AccountChannel\",\"account_id\":1}"}
|
||||
// {"type":"confirm_subscribe","identifier":"..."}
|
||||
// {"type":"confirm_unsubscribe","identifier":"..."}
|
||||
// {"type":"ping","message":"2026-05-23T10:00:00Z"}
|
||||
// {"type":"reject_subscribe","identifier":"...","reason":"..."}
|
||||
|
||||
// CommandType — client→server action types
|
||||
type CommandType string
|
||||
|
||||
const (
|
||||
CommandSubscribe CommandType = "subscribe"
|
||||
CommandUnsubscribe CommandType = "unsubscribe"
|
||||
CommandPing CommandType = "ping"
|
||||
)
|
||||
|
||||
// ServerMessageType — server→client message types
|
||||
type ServerMessageType string
|
||||
|
||||
const (
|
||||
// ServerEvent pushes a real-time event to subscribed clients
|
||||
ServerEvent ServerMessageType = "event"
|
||||
// ServerConfirmSubscribe acknowledges a successful subscription
|
||||
ServerConfirmSubscribe ServerMessageType = "confirm_subscribe"
|
||||
// ServerConfirmUnsubscribe acknowledges a successful unsubscribe
|
||||
ServerConfirmUnsubscribe ServerMessageType = "confirm_unsubscribe"
|
||||
// ServerRejectSubscribe rejects a subscription attempt
|
||||
ServerRejectSubscribe ServerMessageType = "reject_subscribe"
|
||||
// ServerPing is a heartbeat response
|
||||
ServerPing ServerMessageType = "ping"
|
||||
// ServerWelcome is sent immediately upon connection
|
||||
ServerWelcome ServerMessageType = "welcome"
|
||||
// ServerDisconnect is sent before closing the connection
|
||||
ServerDisconnect ServerMessageType = "disconnect"
|
||||
)
|
||||
|
||||
// --- Client → Server Frames ---
|
||||
|
||||
// CommandFrame is the frame clients send to the server.
|
||||
// Mirrors Chatwoot ActionCable's command structure.
|
||||
type CommandFrame struct {
|
||||
Command CommandType `json:"command"`
|
||||
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
|
||||
ConversationID uint `json:"conversation_id,omitempty"` // required for ConversationChannel
|
||||
}
|
||||
|
||||
// Channel name constants (ActionCable naming style)
|
||||
const (
|
||||
ChannelAccount = "AccountChannel"
|
||||
ChannelConversation = "ConversationChannel"
|
||||
)
|
||||
|
||||
// --- Server → Client Frames ---
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// ConfirmFrame acknowledges a subscribe/unsubscribe command.
|
||||
type ConfirmFrame struct {
|
||||
Type ServerMessageType `json:"type"`
|
||||
Identifier string `json:"identifier"`
|
||||
}
|
||||
|
||||
// RejectFrame rejects a subscribe command with a reason.
|
||||
type RejectFrame struct {
|
||||
Type ServerMessageType `json:"type"`
|
||||
Identifier string `json:"identifier"`
|
||||
Reason string `json:"reason"`
|
||||
}
|
||||
|
||||
// PingFrame is a heartbeat pong response.
|
||||
type PingFrame struct {
|
||||
Type ServerMessageType `json:"type"`
|
||||
Message string `json:"message"` // timestamp string
|
||||
}
|
||||
|
||||
// WelcomeFrame is sent upon successful WebSocket connection.
|
||||
type WelcomeFrame struct {
|
||||
Type ServerMessageType `json:"type"`
|
||||
}
|
||||
|
||||
// DisconnectFrame is sent before closing a connection.
|
||||
type DisconnectFrame struct {
|
||||
Type ServerMessageType `json:"type"`
|
||||
Reason string `json:"reason"`
|
||||
Reconnect bool `json:"reconnect"`
|
||||
}
|
||||
|
||||
// --- Real-time Event Type Constants ---
|
||||
// These match the Watermill PubSub topic names and channel.EventType values.
|
||||
|
||||
const (
|
||||
// Message events
|
||||
EventMessageCreated = "message.created"
|
||||
EventMessageUpdated = "message.updated"
|
||||
EventMessageDeleted = "message.deleted"
|
||||
|
||||
// Conversation events
|
||||
EventConversationCreated = "conversation.created"
|
||||
EventConversationUpdated = "conversation.updated"
|
||||
EventConversationResolved = "conversation.resolved"
|
||||
EventConversationOpened = "conversation.opened"
|
||||
EventConversationAssigned = "conversation.assigned"
|
||||
EventConversationUnassigned = "conversation.unassigned"
|
||||
|
||||
// Contact events
|
||||
EventContactCreated = "contact.created"
|
||||
EventContactUpdated = "contact.updated"
|
||||
EventContactDeleted = "contact.deleted"
|
||||
|
||||
// Agent/typing events
|
||||
EventAgentTypingOn = "agent.typing_on"
|
||||
EventAgentTypingOff = "agent.typing_off"
|
||||
EventAgentOnline = "agent.online"
|
||||
EventAgentOffline = "agent.offline"
|
||||
|
||||
// Inbox events
|
||||
EventInboxCreated = "inbox.created"
|
||||
EventInboxUpdated = "inbox.updated"
|
||||
EventInboxDeleted = "inbox.deleted"
|
||||
|
||||
// System notification event
|
||||
EventSystemNotification = "system.notification"
|
||||
|
||||
// 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 ---
|
||||
|
||||
const (
|
||||
// PingInterval is how often the server sends ping frames to detect dead connections.
|
||||
PingInterval = 30 // seconds
|
||||
)
|
||||
@@ -0,0 +1,382 @@
|
||||
package ws
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/ThreeDotsLabs/watermill"
|
||||
"github.com/ThreeDotsLabs/watermill-redisstream/pkg/redisstream"
|
||||
"github.com/ThreeDotsLabs/watermill/message"
|
||||
"github.com/redis/go-redis/v9"
|
||||
|
||||
"github.com/gochat/gochat/pkg/logger"
|
||||
)
|
||||
|
||||
// Subscriber bridges Watermill PubSub events to WebSocket push notifications.
|
||||
// It subscribes to relevant Watermill topics (message, conversation, contact events)
|
||||
// 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).
|
||||
//
|
||||
// Reference: P2E §3 — Real-time communication via Redis Pub/Sub
|
||||
type Subscriber struct {
|
||||
hub *Hub
|
||||
subscriber *redisstream.Subscriber
|
||||
router *message.Router
|
||||
redisClient redis.UniversalClient
|
||||
}
|
||||
|
||||
// NewSubscriber creates a PubSub-to-WebSocket bridge subscriber.
|
||||
func NewSubscriber(hub *Hub, redisClient redis.UniversalClient) (*Subscriber, error) {
|
||||
loggerAdapter := &watermillZapAdapter{}
|
||||
|
||||
wmSubscriber, err := redisstream.NewSubscriber(
|
||||
redisstream.SubscriberConfig{
|
||||
Client: redisClient,
|
||||
ConsumerGroup: "ws-gateway-" + watermill.NewUUID(),
|
||||
},
|
||||
loggerAdapter,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create watermill subscriber: %w", err)
|
||||
}
|
||||
|
||||
router, err := message.NewRouter(message.RouterConfig{}, loggerAdapter)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create watermill router: %w", err)
|
||||
}
|
||||
|
||||
s := &Subscriber{
|
||||
hub: hub,
|
||||
subscriber: wmSubscriber,
|
||||
router: router,
|
||||
redisClient: redisClient,
|
||||
}
|
||||
|
||||
// Register handler functions for each event topic
|
||||
s.registerHandlers()
|
||||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// registerHandlers sets up Watermill router handlers that forward PubSub events
|
||||
// to WebSocket clients via the Hub.
|
||||
func (s *Subscriber) registerHandlers() {
|
||||
// --- Message events → Account room + Conversation room ---
|
||||
// Using AddNoPublisherHandler because these handlers push directly to WS, not to an output topic.
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-message-created-handler",
|
||||
"gochat.message.created",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventMessageCreated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-message-updated-handler",
|
||||
"gochat.message.updated",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventMessageUpdated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-message-deleted-handler",
|
||||
"gochat.message.deleted",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventMessageDeleted),
|
||||
)
|
||||
|
||||
// --- Conversation events → Account room ---
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-conversation-created-handler",
|
||||
"gochat.conversation.created",
|
||||
s.subscriber,
|
||||
s.forwardToAccount(EventConversationCreated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-conversation-updated-handler",
|
||||
"gochat.conversation.updated",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventConversationUpdated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-conversation-resolved-handler",
|
||||
"gochat.conversation.resolved",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventConversationResolved),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-conversation-assigned-handler",
|
||||
"gochat.conversation.assigned",
|
||||
s.subscriber,
|
||||
s.forwardToAccountAndConversation(EventConversationAssigned),
|
||||
)
|
||||
|
||||
// --- Contact events → Account room ---
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-contact-created-handler",
|
||||
"gochat.contact.created",
|
||||
s.subscriber,
|
||||
s.forwardToAccount(EventContactCreated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-contact-updated-handler",
|
||||
"gochat.contact.updated",
|
||||
s.subscriber,
|
||||
s.forwardToAccount(EventContactUpdated),
|
||||
)
|
||||
|
||||
// --- Agent typing events → Conversation room ---
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-agent-typing-on-handler",
|
||||
"gochat.agent.typing_on",
|
||||
s.subscriber,
|
||||
s.forwardToConversation(EventAgentTypingOn),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-agent-typing-off-handler",
|
||||
"gochat.agent.typing_off",
|
||||
s.subscriber,
|
||||
s.forwardToConversation(EventAgentTypingOff),
|
||||
)
|
||||
|
||||
// --- Inbox events → Account room ---
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-inbox-created-handler",
|
||||
"gochat.inbox.created",
|
||||
s.subscriber,
|
||||
s.forwardToAccount(EventInboxCreated),
|
||||
)
|
||||
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-inbox-updated-handler",
|
||||
"gochat.inbox.updated",
|
||||
s.subscriber,
|
||||
s.forwardToAccount(EventInboxUpdated),
|
||||
)
|
||||
|
||||
// --- Notification events → Account room (P4 M8 Notification+Webhook) ---
|
||||
s.router.AddNoPublisherHandler(
|
||||
"ws-notification-created-handler",
|
||||
"gochat.notification.created",
|
||||
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
|
||||
// the account room and the conversation room (if the payload contains conversation_id).
|
||||
// Used for message events and conversation update events that are relevant at both levels.
|
||||
func (s *Subscriber) forwardToAccountAndConversation(eventType string) func(msg *message.Message) error {
|
||||
return func(msg *message.Message) error {
|
||||
payload, err := extractPayload(msg)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
|
||||
return nil // don't retry — bad payload
|
||||
}
|
||||
|
||||
accountID := payload.AccountID
|
||||
if accountID == 0 {
|
||||
logger.L().Warnf("ws: event %s missing account_id", eventType)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Build the event frame
|
||||
frame := EventFrame{
|
||||
Type: ServerEvent,
|
||||
Event: eventType,
|
||||
Payload: payload.Data,
|
||||
}
|
||||
frameData, err := json.Marshal(frame)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to marshal event frame for %s: %v", eventType, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Push to account room
|
||||
s.hub.SendToAccount(accountID, frameData)
|
||||
|
||||
// Also push to conversation room if conversation_id is present
|
||||
if payload.ConversationID > 0 {
|
||||
s.hub.SendToAccountConversation(accountID, payload.ConversationID, frameData)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// forwardToAccount creates a handler that pushes an event only to the account room.
|
||||
// Used for conversation/contact/inbox creation events that don't target a specific conversation.
|
||||
func (s *Subscriber) forwardToAccount(eventType string) func(msg *message.Message) error {
|
||||
return func(msg *message.Message) error {
|
||||
payload, err := extractPayload(msg)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
accountID := payload.AccountID
|
||||
if accountID == 0 {
|
||||
logger.L().Warnf("ws: event %s missing account_id", eventType)
|
||||
return nil
|
||||
}
|
||||
|
||||
frame := EventFrame{
|
||||
Type: ServerEvent,
|
||||
Event: eventType,
|
||||
Payload: payload.Data,
|
||||
}
|
||||
frameData, err := json.Marshal(frame)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to marshal event frame for %s: %v", eventType, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
s.hub.SendToAccount(accountID, frameData)
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// forwardToConversation creates a handler that pushes an event only to the conversation room.
|
||||
// Used for typing events that only need to reach clients viewing that conversation.
|
||||
func (s *Subscriber) forwardToConversation(eventType string) func(msg *message.Message) error {
|
||||
return func(msg *message.Message) error {
|
||||
payload, err := extractPayload(msg)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to extract payload for %s: %v", eventType, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
accountID := payload.AccountID
|
||||
conversationID := payload.ConversationID
|
||||
if accountID == 0 || conversationID == 0 {
|
||||
logger.L().Warnf("ws: event %s missing account_id or conversation_id", eventType)
|
||||
return nil
|
||||
}
|
||||
|
||||
frame := EventFrame{
|
||||
Type: ServerEvent,
|
||||
Event: eventType,
|
||||
Payload: payload.Data,
|
||||
}
|
||||
frameData, err := json.Marshal(frame)
|
||||
if err != nil {
|
||||
logger.L().Errorf("ws: failed to marshal event frame for %s: %v", eventType, err)
|
||||
return nil
|
||||
}
|
||||
|
||||
s.hub.SendToAccountConversation(accountID, conversationID, frameData)
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// --- Event Payload Extraction ---
|
||||
|
||||
// EventPayload represents the common structure of PubSub event messages.
|
||||
// Watermill message payloads are JSON-encoded with account_id, conversation_id,
|
||||
// and an arbitrary data field.
|
||||
type EventPayload struct {
|
||||
AccountID uint `json:"account_id"`
|
||||
ConversationID uint `json:"conversation_id,omitempty"`
|
||||
Data interface{} `json:"data"`
|
||||
}
|
||||
|
||||
// extractPayload parses a Watermill message payload into an EventPayload.
|
||||
func extractPayload(msg *message.Message) (*EventPayload, error) {
|
||||
var payload EventPayload
|
||||
if err := json.Unmarshal(msg.Payload, &payload); err != nil {
|
||||
return nil, fmt.Errorf("json unmarshal failed: %w", err)
|
||||
}
|
||||
return &payload, nil
|
||||
}
|
||||
|
||||
// --- Watermill Zap Logger Adapter ---
|
||||
// Reuses the same pattern as pubsub/event_bus.go's zapLoggerAdapter.
|
||||
|
||||
type watermillZapAdapter struct{}
|
||||
|
||||
func (l *watermillZapAdapter) Error(msg string, err error, fields watermill.LogFields) {
|
||||
if err != nil {
|
||||
logger.L().Errorf("watermill-ws: %s err=%v fields=%v", msg, err, fields)
|
||||
} else {
|
||||
logger.L().Errorf("watermill-ws: %s fields=%v", msg, fields)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *watermillZapAdapter) Info(msg string, fields watermill.LogFields) {
|
||||
logger.L().Infof("watermill-ws: %s fields=%v", msg, fields)
|
||||
}
|
||||
|
||||
func (l *watermillZapAdapter) Debug(msg string, fields watermill.LogFields) {
|
||||
// Skip debug logs to reduce noise
|
||||
}
|
||||
|
||||
func (l *watermillZapAdapter) Trace(msg string, fields watermill.LogFields) {
|
||||
// Skip trace logs
|
||||
}
|
||||
|
||||
func (l *watermillZapAdapter) With(fields watermill.LogFields) watermill.LoggerAdapter {
|
||||
return l
|
||||
}
|
||||
|
||||
// --- Lifecycle Management ---
|
||||
|
||||
// Run starts the Watermill router (blocking). Should be called in a goroutine.
|
||||
func (s *Subscriber) Run(ctx context.Context) error {
|
||||
logger.L().Info("ws: starting PubSub subscriber router")
|
||||
return s.router.Run(ctx)
|
||||
}
|
||||
|
||||
// Close gracefully shuts down the Watermill router and subscriber.
|
||||
func (s *Subscriber) Close() error {
|
||||
logger.L().Info("ws: shutting down PubSub subscriber")
|
||||
if err := s.router.Close(); err != nil {
|
||||
logger.L().Errorf("ws: router close error: %v", err)
|
||||
return err
|
||||
}
|
||||
if err := s.subscriber.Close(); err != nil {
|
||||
if !strings.Contains(err.Error(), "already closed") {
|
||||
logger.L().Errorf("ws: subscriber close error: %v", err)
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Running returns whether the router is currently running.
|
||||
func (s *Subscriber) Running() bool {
|
||||
return s.router.IsRunning()
|
||||
}
|
||||
Reference in New Issue
Block a user