diff --git a/README.md b/README.md index 126881f..46b569a 100644 --- a/README.md +++ b/README.md @@ -136,9 +136,9 @@ Keeps the connection alive. #### `broadcast` -Sent by the server to clients when a message is published to a channel they are subscribed to. +Sent by the server to clients when a message is published to a channel they are subscribed to. The `seq` field is an incrementing integer for each subscription, which can be used to detect missed events. -**Params**: `{"id": "msg-123", "createTime": "2023-01-01T12:00:00Z", "channel": "channel-name", "event": "event-name", "payload": {"key": "value"}}` +**Params**: `{"id": "msg-123", "seq": 1, "createTime": "2023-01-01T12:00:00Z", "channel": "channel-name", "event": "event-name", "payload": {"key": "value"}}` ## REST API diff --git a/internal/broadcaster/connection.go b/internal/broadcaster/connection.go index d8536d2..387624c 100644 --- a/internal/broadcaster/connection.go +++ b/internal/broadcaster/connection.go @@ -12,9 +12,18 @@ type Connection struct { Send chan Message mu sync.RWMutex + Seq uint64 authentication *auth.Authentication } +func (c *Connection) NextSeq() uint64 { + c.mu.Lock() + defer c.mu.Unlock() + + c.Seq++ + return c.Seq +} + func (c *Connection) SetAuthentication(auth *auth.Authentication) { c.mu.Lock() defer c.mu.Unlock() diff --git a/internal/broadcaster/message.go b/internal/broadcaster/message.go index 79282af..0bd65b6 100644 --- a/internal/broadcaster/message.go +++ b/internal/broadcaster/message.go @@ -4,6 +4,7 @@ import "time" type Message struct { Id string `json:"id"` + Seq uint64 `json:"seq"` CreateTime time.Time `json:"createTime"` Channel string `json:"channel"` Event string `json:"event"` diff --git a/internal/broadcaster/registry.go b/internal/broadcaster/registry.go index 2c15276..2f7b387 100644 --- a/internal/broadcaster/registry.go +++ b/internal/broadcaster/registry.go @@ -69,8 +69,11 @@ func (r *InMemoryRegistry) Broadcast(message Message) { var staleConnectionIds []string for _, connection := range connections { + msg := message + msg.Seq = connection.NextSeq() + select { - case connection.Send <- message: + case connection.Send <- msg: default: r.logger.Warn("connection send channel is full, closing connection", zap.String("connectionId", connection.Id)) diff --git a/internal/server/websocket.go b/internal/server/websocket.go index b69832c..733c13e 100644 --- a/internal/server/websocket.go +++ b/internal/server/websocket.go @@ -52,6 +52,7 @@ func (s *WebSocketServer) Register(router *mux.Router) { broadcasterConn := &broadcaster.Connection{ Id: connectionId, Send: broascasterChannel, + Seq: 0, } s.registry.Connect(broadcasterConn)