package pubsub import ( "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" ) // Reference: P2E §3 — Real-time communication via Redis Pub/Sub // Using Watermill as the unified message bus, replacing Chatwoot's Wisper + Sidekiq pattern. // // Architecture mapping: // Chatwoot Wisper (sync broadcast) → Watermill Router (sync handler) // Chatwoot Sidekiq (async workers) → Watermill Redis Streams (persistent async) // Chatwoot ActionCable (WebSocket) → Watermill subscriber → WebSocket gateway // // Watermill advantages over asynq: // - Multiple pub/sub backends (Redis Streams, Kafka, AMQP, SQL, Google Cloud Pub/Sub) // - Built-in Router for event routing (like Chatwoot's Wisper event names) // - CQRS component for command/query separation // - At-least-once delivery with ack/nack semantics // - Consumer groups for horizontal scaling // --- Topic Constants --- // Topics follow Chatwoot's event naming convention with GoChat-specific prefixes. // Format: gochat.{domain}.{event_type} const ( // Message topics (ref: Chatwoot MessageCreated/MessageUpdated events) TopicMessageCreated = "gochat.message.created" TopicMessageUpdated = "gochat.message.updated" TopicMessageDeleted = "gochat.message.deleted" // Conversation topics (ref: Chatwoot ConversationStatusChanged events) TopicConversationCreated = "gochat.conversation.created" TopicConversationUpdated = "gochat.conversation.updated" TopicConversationResolved = "gochat.conversation.resolved" TopicConversationAssigned = "gochat.conversation.assigned" // 1:1 Chatwoot: additional conversation events TopicConversationStatusChanged = "gochat.conversation.status_changed" TopicConversationContactChanged = "gochat.conversation.contact_changed" TopicConversationRead = "gochat.conversation.read" TopicConversationUnreadCountChanged = "gochat.conversation.unread_count_changed" TopicConversationMentioned = "gochat.conversation.mentioned" TopicAssigneeChanged = "gochat.conversation.assignee_changed" TopicTeamChanged = "gochat.conversation.team_changed" TopicFirstReplyCreated = "gochat.message.first_reply_created" // Contact topics (ref: Chatwoot ContactCreated/ContactUpdated events) TopicContactCreated = "gochat.contact.created" TopicContactUpdated = "gochat.contact.updated" // 1:1 Chatwoot: additional contact events TopicContactMerged = "gochat.contact.merged" TopicContactDeleted = "gochat.contact.deleted" // Channel event topics (ref: Chatwoot webhook events per channel) TopicChannelEvent = "gochat.channel.event.%s" // %s = channel_type (telegram, web_widget, etc.) // Inbox topics TopicInboxCreated = "gochat.inbox.created" TopicInboxUpdated = "gochat.inbox.updated" TopicInboxDeleted = "gochat.inbox.deleted" // Agent topics (ref: Chatwoot Agent events) TopicAgentTypingOn = "gochat.agent.typing_on" TopicAgentTypingOff = "gochat.agent.typing_off" // System topics TopicWebhookReceived = "gochat.webhook.received" TopicSystemNotification = "gochat.system.notification" TopicNotificationCreated = "gochat.notification.created" TopicNotificationDeleted = "gochat.notification.deleted" // Account topics TopicAccountCacheInvalidated = "gochat.account.cache_invalidated" // AccountUser topics — 1:1 Chatwoot: after_create_commit/after_destroy callbacks TopicAccountUserCreated = "gochat.account_user.created" TopicAccountUserUpdated = "gochat.account_user.updated" TopicAccountUserDestroyed = "gochat.account_user.destroyed" ) // FormatTopic creates a topic string with a dynamic suffix. // Example: FormatTopic(TopicChannelEvent, "telegram") → "gochat.channel.event.telegram" func FormatTopic(template, suffix string) string { result := make([]byte, 0, len(template)+len(suffix)) i := 0 for i < len(template) { if i+1 < len(template) && template[i] == '%' && template[i+1] == 's' { result = append(result, suffix...) i += 2 } else { result = append(result, template[i]) i++ } } return string(result) } // --- Watermill Adapter --- // EventBus wraps Watermill publisher for application-level event publishing. // This is the GoChat equivalent of Chatwoot's Wisper.broadcast(). type EventBus struct { publisher message.Publisher } // NewEventBus creates an EventBus backed by Redis Streams via Watermill. func NewEventBus(redisClient redis.UniversalClient) (*EventBus, error) { loggerAdapter := &zapLoggerAdapter{} publisher, err := redisstream.NewPublisher( redisstream.PublisherConfig{ Client: redisClient, }, loggerAdapter, ) if err != nil { return nil, err } return &EventBus{publisher: publisher}, nil } // Publish publishes an event to a topic. // Payload should be JSON-encoded event data. func (bus *EventBus) Publish(topic string, payload []byte) error { msg := message.NewMessage(watermill.NewUUID(), payload) msg.Metadata.Set("topic", topic) msg.Metadata.Set("source", "gochat") if err := bus.publisher.Publish(topic, msg); err != nil { logger.L().Errorf("failed to publish event to topic %s: %v", topic, err) return err } logger.L().Infof("published event to topic %s (msg_id=%s)", topic, msg.UUID) return nil } // Close shuts down the publisher. func (bus *EventBus) Close() error { return bus.publisher.Close() } // --- Subscriber Factory --- // NewSubscriber creates a Watermill Redis Streams subscriber for a consumer group. // Each GoChat service instance runs as a separate consumer within the group, // enabling horizontal scaling (same pattern as Chatwoot's Sidekiq workers). func NewSubscriber(redisClient redis.UniversalClient, consumerGroup string) (*redisstream.Subscriber, error) { loggerAdapter := &zapLoggerAdapter{} subscriber, err := redisstream.NewSubscriber( redisstream.SubscriberConfig{ Client: redisClient, ConsumerGroup: consumerGroup, BlockTime: redisstream.DefaultBlockTime, ClaimInterval: redisstream.DefaultClaimInterval, ClaimBatchSize: redisstream.DefaultClaimBatchSize, MaxIdleTime: redisstream.DefaultMaxIdleTime, CheckConsumersInterval: redisstream.DefaultCheckConsumersInterval, ConsumerTimeout: redisstream.DefaultConsumerTimeout, }, loggerAdapter, ) if err != nil { return nil, err } return subscriber, nil } // --- Zap Logger Adapter --- // Bridges GoChat's zap.SugaredLogger to Watermill's LoggerAdapter interface. type zapLoggerAdapter struct{} func (l *zapLoggerAdapter) Error(msg string, err error, fields watermill.LogFields) { if err != nil { logger.L().Errorf("watermill: %s err=%v fields=%v", msg, err, fields) } else { logger.L().Errorf("watermill: %s fields=%v", msg, fields) } } func (l *zapLoggerAdapter) Info(msg string, fields watermill.LogFields) { logger.L().Infof("watermill: %s fields=%v", msg, fields) } func (l *zapLoggerAdapter) Debug(msg string, fields watermill.LogFields) { // Skip debug logs to reduce noise in production } func (l *zapLoggerAdapter) Trace(msg string, fields watermill.LogFields) { // Skip trace logs } func (l *zapLoggerAdapter) With(fields watermill.LogFields) watermill.LoggerAdapter { return l // No field enrichment needed — zap handles structured logging via .L() }