From 1db4408244fe75b39c2b4e1bbb8e681679dd7dae Mon Sep 17 00:00:00 2001 From: prakash-dev-code Date: Wed, 26 Aug 2026 20:05:41 +0000 Subject: [PATCH] fix(websocket): deliver events to every subscriber of an instance The websocket producer kept a single connection per instance in `clients map[string]*websocket.Conn`. Registering a second subscriber silently replaced the first, and `RemoveClient` deleted the whole instance entry when either one disconnected. Consuming events in two browser tabs therefore delivered to only the most recent one, and closing that tab cut delivery for the remaining subscriber too. Changes: - `clients` becomes `map[string][]*client`. Deregistration now matches the specific connection instead of removing the instance entry, so sibling subscribers keep receiving events. - Each connection carries its own write mutex. Events are dispatched from independent goroutines (the `go CallWebhook` calls in the whatsmeow event handler) and gorilla/websocket permits only one concurrent writer per connection, so two simultaneous events could interleave writes on the same socket. - `Produce` snapshots the subscriber list under the read lock and performs writes outside it, so a slow or half-open socket no longer blocks deliveries for other instances. - Failed writes prune the dead connection, and `Produce` now returns nil. `sendToQueueOrWebhook` aborts the remaining producers when one returns an error, so a closed browser tab previously suppressed RabbitMQ, NATS and webhook delivery for that same event. Verified with two concurrent subscribers on one instance: both receive identical event streams, and disconnecting one leaves the other unaffected. --- pkg/events/websocket/websocket_producer.go | 171 ++++++++++++++------- 1 file changed, 112 insertions(+), 59 deletions(-) diff --git a/pkg/events/websocket/websocket_producer.go b/pkg/events/websocket/websocket_producer.go index f28ff9aa..bf37a271 100644 --- a/pkg/events/websocket/websocket_producer.go +++ b/pkg/events/websocket/websocket_producer.go @@ -14,92 +14,131 @@ var upgrader = websocket.Upgrader{ ReadBufferSize: 1024, WriteBufferSize: 1024, CheckOrigin: func(r *http.Request) bool { - logger.LogInfo("Verificando origem da conexão WebSocket") return true }, } +// client wraps a connection together with its own write mutex. +// +// gorilla/websocket permits only one concurrent writer per connection, and +// events are dispatched from independent goroutines (see the `go CallWebhook` +// calls in the whatsmeow event handler), so two events arriving at once used to +// race on the same socket. Serialising per connection removes that. +type client struct { + conn *websocket.Conn + mu sync.Mutex +} + +func (c *client) writeJSON(v interface{}) error { + c.mu.Lock() + defer c.mu.Unlock() + return c.conn.WriteJSON(v) +} + type websocketProducer struct { - clients map[string]*websocket.Conn // conexões específicas por instância - broadcast []*websocket.Conn // conexões que recebem todos os eventos + // clients holds every connection subscribed to a given instance. This is a + // slice, not a single connection: previously a second subscriber replaced the + // first, so opening the UI in two tabs silently cut delivery to the older one, + // and either tab disconnecting removed the whole instance entry. + clients map[string][]*client + broadcast []*client clientsMux sync.RWMutex loggerWrapper *logger_wrapper.LoggerManager } func NewWebsocketProducer(loggerWrapper *logger_wrapper.LoggerManager) *websocketProducer { return &websocketProducer{ - clients: make(map[string]*websocket.Conn), - broadcast: make([]*websocket.Conn, 0), - clientsMux: sync.RWMutex{}, + clients: make(map[string][]*client), + broadcast: make([]*client, 0), loggerWrapper: loggerWrapper, } } -// ServeWs lida com as requisições de upgrade para websocket +// ServeWs upgrades an HTTP request and registers it for the lifetime of the +// connection. An empty instanceId subscribes to every instance's events. func ServeWs(w http.ResponseWriter, r *http.Request, instanceId string, producer *websocketProducer) { - logger.LogInfo("Iniciando upgrade da conexão WebSocket") conn, err := upgrader.Upgrade(w, r, nil) if err != nil { logger.LogError("Erro ao fazer upgrade da conexão websocket: %v", err) return } - logger.LogInfo("Conexão WebSocket estabelecida com sucesso") - + c := &client{conn: conn} if instanceId == "" { - producer.AddBroadcastClient(conn) + producer.addBroadcastClient(c) } else { - producer.AddClient(instanceId, conn) + producer.addClient(instanceId, c) } - // Goroutine para limpar conexão quando fechada + // The read loop exists purely to notice a closed socket; inbound frames are + // ignored. Deregistration is keyed on this exact connection so sibling + // subscribers of the same instance keep receiving events. go func() { + defer func() { + if instanceId == "" { + producer.removeBroadcastClient(c) + } else { + producer.removeClient(instanceId, c) + } + _ = conn.Close() + }() for { - _, _, err := conn.ReadMessage() - if err != nil { - if instanceId == "" { - producer.RemoveBroadcastClient(conn) - } else { - producer.RemoveClient(instanceId) - } - conn.Close() - break + if _, _, err := conn.ReadMessage(); err != nil { + return } } }() } -func (p *websocketProducer) AddBroadcastClient(conn *websocket.Conn) { +func (p *websocketProducer) addBroadcastClient(c *client) { p.clientsMux.Lock() - defer p.clientsMux.Unlock() - p.broadcast = append(p.broadcast, conn) - logger.LogInfo("Cliente broadcast websocket adicionado") + p.broadcast = append(p.broadcast, c) + n := len(p.broadcast) + p.clientsMux.Unlock() + logger.LogInfo("Cliente broadcast websocket adicionado (total: %d)", n) } -func (p *websocketProducer) RemoveBroadcastClient(conn *websocket.Conn) { +func (p *websocketProducer) removeBroadcastClient(c *client) { p.clientsMux.Lock() - defer p.clientsMux.Unlock() - for i, c := range p.broadcast { - if c == conn { - p.broadcast = append(p.broadcast[:i], p.broadcast[i+1:]...) - break - } - } - logger.LogInfo("Cliente broadcast websocket removido") + p.broadcast = drop(p.broadcast, c) + n := len(p.broadcast) + p.clientsMux.Unlock() + logger.LogInfo("Cliente broadcast websocket removido (total: %d)", n) } -func (p *websocketProducer) AddClient(instanceID string, conn *websocket.Conn) { +func (p *websocketProducer) addClient(instanceID string, c *client) { p.clientsMux.Lock() - defer p.clientsMux.Unlock() - p.clients[instanceID] = conn - p.loggerWrapper.GetLogger(instanceID).LogInfo("Cliente websocket adicionado para instância: %s", instanceID) + p.clients[instanceID] = append(p.clients[instanceID], c) + n := len(p.clients[instanceID]) + p.clientsMux.Unlock() + p.loggerWrapper.GetLogger(instanceID).LogInfo( + "Cliente websocket adicionado para instância: %s (total: %d)", instanceID, n) } -func (p *websocketProducer) RemoveClient(instanceID string) { +func (p *websocketProducer) removeClient(instanceID string, c *client) { p.clientsMux.Lock() - defer p.clientsMux.Unlock() - delete(p.clients, instanceID) - p.loggerWrapper.GetLogger(instanceID).LogInfo("Cliente websocket removido para instância: %s", instanceID) + remaining := drop(p.clients[instanceID], c) + if len(remaining) == 0 { + delete(p.clients, instanceID) + } else { + p.clients[instanceID] = remaining + } + p.clientsMux.Unlock() + p.loggerWrapper.GetLogger(instanceID).LogInfo( + "Cliente websocket removido para instância: %s (restantes: %d)", instanceID, len(remaining)) +} + +// drop returns a new slice without target. It allocates rather than filtering in +// place, so a snapshot taken by a concurrent Produce can never observe a mutated +// backing array. +func drop(list []*client, target *client) []*client { + out := make([]*client, 0, len(list)) + for _, c := range list { + if c != target { + out = append(out, c) + } + } + return out } func (p *websocketProducer) Produce(queueName string, payload []byte, instanceID string, _ string) error { @@ -108,29 +147,43 @@ func (p *websocketProducer) Produce(queueName string, payload []byte, instanceID "payload": string(payload), } + // Snapshot under the lock, then write outside it: a slow or half-open socket + // must not block other instances' deliveries. p.clientsMux.RLock() - defer p.clientsMux.RUnlock() - - // Envia para cliente específico da instância - if client, exists := p.clients[instanceID]; exists { - err := client.WriteJSON(message) - if err != nil { - p.loggerWrapper.GetLogger(instanceID).LogError("Erro ao enviar mensagem websocket para %s: %v", instanceID, err) - // Não remove o cliente aqui pois estamos com o RLock - return err + targets := append([]*client(nil), p.clients[instanceID]...) + targets = append(targets, p.broadcast...) + p.clientsMux.RUnlock() + + var failed []*client + for _, c := range targets { + if err := c.writeJSON(message); err != nil { + p.loggerWrapper.GetLogger(instanceID).LogError( + "Erro ao enviar mensagem websocket para %s: %v", instanceID, err) + failed = append(failed, c) } - p.loggerWrapper.GetLogger(instanceID).LogInfo("Mensagem websocket enviada com sucesso para instância %s na fila %s", instanceID, queueName) } - // Envia para todos os clientes broadcast - for _, conn := range p.broadcast { - err := conn.WriteJSON(message) - if err != nil { - p.loggerWrapper.GetLogger(instanceID).LogError("Erro ao enviar mensagem broadcast websocket: %v", err) - continue + if len(failed) > 0 { + p.clientsMux.Lock() + for _, c := range failed { + if remaining := drop(p.clients[instanceID], c); len(remaining) == 0 { + delete(p.clients, instanceID) + } else { + p.clients[instanceID] = remaining + } + p.broadcast = drop(p.broadcast, c) } + p.clientsMux.Unlock() + } + + if n := len(targets) - len(failed); n > 0 { + p.loggerWrapper.GetLogger(instanceID).LogInfo( + "Mensagem websocket enviada para %d cliente(s) da instância %s na fila %s", n, instanceID, queueName) } + // Deliberately nil even when a write failed. sendToQueueOrWebhook aborts the + // remaining producers on error, so reporting a dead browser tab here would + // silently suppress RabbitMQ/NATS/webhook delivery for the same event. return nil }