Files
gochat/backend/internal/ws/event_publisher.go
T
rogee aeddedf2a3 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.
2026-07-07 14:44:12 +08:00

203 lines
7.2 KiB
Go

// Package ws provides the EventPublisher that unifies real-time event delivery
// across WebSocket Hub and SSE fallback channels. Services call PublishEvent
// after mutations (CreateConversation, SendMessage, etc.) to push events to
// all connected clients, mirroring Chatwoot's ActionCable broadcast pattern.
package ws
import (
"context"
"encoding/json"
"fmt"
applogger "github.com/gochat/gochat/pkg/logger"
)
// EventPublisher routes business events to both WebSocket and SSE clients.
// It is the central point that services call after mutations to push real-time
// updates, matching Chatwoot's pattern where controllers broadcast events
// after create/update/delete actions via Wisper → ActionCable.
//
// Architecture:
//
// Service.CreateX() → EventPublisher.PublishEvent(accountID, eventType, payload)
// → Hub.SendToAccount (WebSocket local delivery)
// → SSERegistry.SendToAccount (SSE local delivery)
// → BroadcastRelay.Publish (Redis Pub/Sub cross-instance delivery)
//
// Reference: Chatwoot uses Wisper (in-process pub/sub) + ActionCable (WebSocket)
// + Redis Pub/Sub for cross-instance. GoChat uses EventPublisher as the unified
// entry point, with Hub and SSERegistry as local delivery targets, and Redis
// Pub/Sub relay for multi-instance fan-out.
type EventPublisher struct {
hub MessageHandler // WebSocket Hub (local delivery)
sse *SSERegistry // SSE registry (local delivery)
relay *BroadcastRelay // Redis Pub/Sub relay (cross-instance delivery)
}
// NewEventPublisher creates an EventPublisher with all delivery targets.
func NewEventPublisher(hub MessageHandler, sse *SSERegistry, relay *BroadcastRelay) *EventPublisher {
return &EventPublisher{
hub: hub,
sse: sse,
relay: relay,
}
}
// NewEventPublisherLocal creates an EventPublisher with only local delivery
// (no Redis relay). Used for development/testing without Redis.
func NewEventPublisherLocal(hub MessageHandler, sse *SSERegistry) *EventPublisher {
return &EventPublisher{
hub: hub,
sse: sse,
}
}
// PublishEvent publishes a real-time event to all delivery targets.
// This is the primary API that services call after business mutations.
//
// Parameters:
// - accountID: the account scope for the event (each client only sees their account's events)
// - eventType: dot-separated event name (e.g. "message.created", "conversation.updated")
// - payload: event data (will be JSON-encoded for transport)
//
// The event is routed to:
// 1. WebSocket Hub → all locally connected WS clients for that account
// 2. SSE Registry → all locally connected SSE clients for that account
// 3. Redis Pub/Sub relay → all other GoChat instances (cross-instance delivery)
//
// Reference: Chatwoot controllers call broadcast_event after mutations,
// which triggers Wisper → ActionCable → Redis Pub/Sub relay.
func (p *EventPublisher) PublishEvent(accountID uint, eventType string, payload interface{}) {
// Build the wire-format message (matching Chatwoot ActionCable event format)
wsMsg := &WSMessage{
Event: eventType,
Data: payload,
AccountID: accountID,
}
data, err := json.Marshal(wsMsg)
if err != nil {
applogger.L().Warnf("event publisher: failed to marshal event %s: %v", eventType, err)
return
}
// 1. Deliver to WebSocket Hub (local clients)
if p.hub != nil {
p.hub.SendToAccount(accountID, data)
}
// 2. Deliver to SSE registry (local clients)
if p.sse != nil {
p.sse.SendToAccount(accountID, SSEEvent{
Type: eventType,
Payload: payload,
})
}
// 3. Publish to Redis Pub/Sub relay (cross-instance delivery)
if p.relay != nil {
room := accountRoomNameHelper(accountID)
if err := p.relay.Publish(context.Background(), room, wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for %s: %v", eventType, err)
}
}
}
// PublishConversationEvent publishes a real-time event scoped to a specific
// conversation within an account. This is used for events that should only
// reach clients viewing that conversation (e.g. typing indicators, message updates).
//
// Reference: Chatwoot ConversationChannel — events broadcast to clients
// subscribed to a specific conversation's channel.
func (p *EventPublisher) PublishConversationEvent(accountID uint, conversationID uint, eventType string, payload interface{}) {
wsMsg := &WSMessage{
Event: eventType,
Data: payload,
AccountID: accountID,
}
data, err := json.Marshal(wsMsg)
if err != nil {
applogger.L().Warnf("event publisher: failed to marshal conversation event %s: %v", eventType, err)
return
}
// 1. Deliver to WebSocket Hub (both account room and conversation room)
if p.hub != nil {
p.hub.SendToAccount(accountID, data)
convRoom := conversationRoomNameHelper(accountID, conversationID)
p.hub.SendToRoom(convRoom, data)
}
// 2. Deliver to SSE registry (conversation-filtered channels)
if p.sse != nil {
p.sse.SendToConversation(accountID, conversationID, SSEEvent{
Type: eventType,
Payload: payload,
})
}
// 3. Publish to Redis Pub/Sub relay (both account and conversation channels)
if p.relay != nil {
room := accountRoomNameHelper(accountID)
if err := p.relay.Publish(context.Background(), room, wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for %s: %v", eventType, err)
}
}
}
// PublishWidgetEvent publishes an event to the account room and the widget
// contact's pubsub_token room. Chatwoot's widget ActionCable connector
// subscribes to RoomChannel with pubsub_token, so widget-visible events must be
// available on that token-scoped room in addition to the dashboard account room.
func (p *EventPublisher) PublishWidgetEvent(accountID uint, pubsubToken string, eventType string, payload interface{}) {
wsMsg := &WSMessage{
Event: eventType,
Data: payload,
AccountID: accountID,
}
data, err := json.Marshal(wsMsg)
if err != nil {
applogger.L().Warnf("event publisher: failed to marshal widget event %s: %v", eventType, err)
return
}
if p.hub != nil {
p.hub.SendToAccount(accountID, data)
if pubsubToken != "" {
p.hub.SendToRoom(pubsubTokenRoomNameHelper(pubsubToken), data)
}
}
if p.sse != nil {
p.sse.SendToAccount(accountID, SSEEvent{Type: eventType, Payload: payload})
}
if p.relay != nil {
room := accountRoomNameHelper(accountID)
if err := p.relay.Publish(context.Background(), room, wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for %s: %v", eventType, err)
}
if pubsubToken != "" {
if err := p.relay.Publish(context.Background(), pubsubTokenRoomNameHelper(pubsubToken), wsMsg); err != nil {
applogger.L().Warnf("event publisher: redis publish failed for widget %s: %v", eventType, err)
}
}
}
}
// accountRoomNameHelper generates the room name for an account channel.
func accountRoomNameHelper(accountID uint) string {
return fmt.Sprintf("account_%d", accountID)
}
// conversationRoomNameHelper generates the room name for a conversation channel.
func conversationRoomNameHelper(accountID uint, conversationID uint) string {
return fmt.Sprintf("account_%d_conversation_%d", accountID, conversationID)
}
func pubsubTokenRoomNameHelper(token string) string {
return fmt.Sprintf("pubsub_token_%s", token)
}