Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
9 changes: 9 additions & 0 deletions internal/broadcaster/connection.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
1 change: 1 addition & 0 deletions internal/broadcaster/message.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
5 changes: 4 additions & 1 deletion internal/broadcaster/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
1 change: 1 addition & 0 deletions internal/server/websocket.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ func (s *WebSocketServer) Register(router *mux.Router) {
broadcasterConn := &broadcaster.Connection{
Id: connectionId,
Send: broascasterChannel,
Seq: 0,
}

s.registry.Connect(broadcasterConn)
Expand Down