Skip to content

Commit

Permalink
Lock around websocket writes
Browse files Browse the repository at this point in the history
  • Loading branch information
ericvolp12 committed Sep 25, 2024
1 parent 25cc0a9 commit 2cd1b61
Show file tree
Hide file tree
Showing 2 changed files with 12 additions and 4 deletions.
6 changes: 3 additions & 3 deletions pkg/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ func (s *Server) HandleSubscribe(c echo.Context) error {

switch msgType {
case websocket.PingMessage:
if err := ws.WriteMessage(websocket.PongMessage, nil); err != nil {
if err := sub.WriteMessage(websocket.PongMessage, nil); err != nil {
log.Error("failed to write pong to websocket", "error", err)
cancel()
return
Expand Down Expand Up @@ -241,15 +241,15 @@ func (s *Server) HandleSubscribe(c echo.Context) error {

// When compression is enabled, the msg is a zstd compressed message
if compress {
if err := ws.WriteMessage(websocket.BinaryMessage, *msg); err != nil {
if err := sub.WriteMessage(websocket.BinaryMessage, *msg); err != nil {
log.Error("failed to write message to websocket", "error", err)
return nil
}
continue
}

// Otherwise, the msg is serialized JSON
if err := ws.WriteMessage(websocket.TextMessage, *msg); err != nil {
if err := sub.WriteMessage(websocket.TextMessage, *msg); err != nil {
log.Error("failed to write message to websocket", "error", err)
return nil
}
Expand Down
10 changes: 9 additions & 1 deletion pkg/server/subscriber.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ type WantedCollections struct {

type Subscriber struct {
ws *websocket.Conn
conLk sync.Mutex
realIP string
lk sync.Mutex
seq int64
Expand Down Expand Up @@ -266,7 +267,14 @@ func (s *Subscriber) UpdateOptions(opts *SubscriberOptions) {

// Terminate sends a close message to the subscriber
func (s *Subscriber) Terminate(reason string) error {
return s.ws.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, reason))
return s.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, reason))
}

func (s *Subscriber) WriteMessage(msgType int, data []byte) error {
s.conLk.Lock()
defer s.conLk.Unlock()

return s.ws.WriteMessage(msgType, data)
}

func (s *Subscriber) SetCursor(cursor *int64) {
Expand Down

0 comments on commit 2cd1b61

Please sign in to comment.