package ws import ( "context" "encoding/json" "fmt" "strings" "github.com/ThreeDotsLabs/watermill" "github.com/ThreeDotsLabs/watermill/message" "github.com/ThreeDotsLabs/watermill-redisstream/pkg/redisstream" "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), ) } // 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() }