MM-34389: Reliable websockets: First commit (#17297)
* MM-34389: Reliable websockets: First commit This is the first commit which makes some basic changes to get it ready for the actual implementation. Changes include: - A config field to conditionally enable it. - Refactoring the WriteMessage along with setting the deadline to a separate method. The basic idea is that the client sends the connection_id and sequence_number either during the handshake (via query params), or during challenge_auth (via added parameters in the map). If the conn_id is empty, then we create a new one and set it. Otherwise, we get the queues from a connection manager (TBD) and attach them to WebConn. ```release-note NONE ``` https://mattermost.atlassian.net/browse/MM-34389 * Incorporate review comments * Trigger CI * removing telemetry Co-authored-by: Mattermod <mattermod@users.noreply.github.com>
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
ca9f4c9ed8
Коммит
4ba0c09fc7
@@ -5,10 +5,17 @@ package api4
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"strconv"
|
||||||
|
|
||||||
"github.com/gorilla/websocket"
|
"github.com/gorilla/websocket"
|
||||||
|
|
||||||
"github.com/mattermost/mattermost-server/v5/model"
|
"github.com/mattermost/mattermost-server/v5/model"
|
||||||
|
"github.com/mattermost/mattermost-server/v5/shared/mlog"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
connectionIDParam = "connection_id"
|
||||||
|
sequenceNumberParam = "sequence_number"
|
||||||
)
|
)
|
||||||
|
|
||||||
func (api *API) InitWebSocket() {
|
func (api *API) InitWebSocket() {
|
||||||
@@ -31,6 +38,51 @@ func connectWebSocket(c *Context, w http.ResponseWriter, r *http.Request) {
|
|||||||
|
|
||||||
wc := c.App.NewWebConn(ws, *c.App.Session(), c.App.T, "")
|
wc := c.App.NewWebConn(ws, *c.App.Session(), c.App.T, "")
|
||||||
|
|
||||||
|
if *c.App.Config().ServiceSettings.EnableReliableWebSockets {
|
||||||
|
connID := r.URL.Query().Get(connectionIDParam)
|
||||||
|
if connID == "" {
|
||||||
|
// If not present, we assume client is not capable yet, or it's a fresh connection.
|
||||||
|
// We just create a new ID.
|
||||||
|
connID = model.NewId()
|
||||||
|
} else {
|
||||||
|
if !model.IsValidId(connID) {
|
||||||
|
mlog.Error("Invalid connection ID", mlog.String("id", connID))
|
||||||
|
wc.WebSocket.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// If present, we check if it's present in the connection manager.
|
||||||
|
// TODO: the connection manager internally should forward the request
|
||||||
|
// to the cluster if it does not have it.
|
||||||
|
//
|
||||||
|
// If the connection is not present, then we assume either timeout,
|
||||||
|
// or server restart. In that case, we set a new one.
|
||||||
|
//
|
||||||
|
// Now we get the sequence number
|
||||||
|
seqVal := r.URL.Query().Get(sequenceNumberParam)
|
||||||
|
if seqVal == "" {
|
||||||
|
// Sequence_number must be sent with connection id.
|
||||||
|
// A client must be either non-compliant or fully compliant.
|
||||||
|
mlog.Error("Sequence number not present in websocket request")
|
||||||
|
wc.WebSocket.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
seq, err := strconv.Atoi(seqVal)
|
||||||
|
if err != nil || seq < 0 {
|
||||||
|
mlog.Error("Invalid sequence number set in query param",
|
||||||
|
mlog.String("query", seqVal),
|
||||||
|
mlog.Err(err))
|
||||||
|
wc.WebSocket.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
wc.Sequence = int64(seq)
|
||||||
|
// Now if there have been past entries to be back-filled, we do it.
|
||||||
|
// First we find the right sequence number point.
|
||||||
|
// We start consuming from dead queue first, and then move to active queue
|
||||||
|
}
|
||||||
|
// In case of fresh connection id, sequence number is already zero.
|
||||||
|
wc.SetConnectionID(connID)
|
||||||
|
}
|
||||||
|
|
||||||
if c.App.Session().UserId != "" {
|
if c.App.Session().UserId != "" {
|
||||||
c.App.HubRegister(wc)
|
c.App.HubRegister(wc)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -49,6 +49,7 @@ type WebConn struct {
|
|||||||
send chan model.WebSocketMessage
|
send chan model.WebSocketMessage
|
||||||
sessionToken atomic.Value
|
sessionToken atomic.Value
|
||||||
session atomic.Value
|
session atomic.Value
|
||||||
|
connectionID atomic.Value
|
||||||
endWritePump chan struct{}
|
endWritePump chan struct{}
|
||||||
pumpFinished chan struct{}
|
pumpFinished chan struct{}
|
||||||
}
|
}
|
||||||
@@ -118,6 +119,11 @@ func (wc *WebConn) SetSessionToken(v string) {
|
|||||||
wc.sessionToken.Store(v)
|
wc.sessionToken.Store(v)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SetConnectionID sets the connection id of the connection.
|
||||||
|
func (wc *WebConn) SetConnectionID(id string) {
|
||||||
|
wc.connectionID.Store(id)
|
||||||
|
}
|
||||||
|
|
||||||
// GetSession returns the session of the connection.
|
// GetSession returns the session of the connection.
|
||||||
func (wc *WebConn) GetSession() *model.Session {
|
func (wc *WebConn) GetSession() *model.Session {
|
||||||
return wc.session.Load().(*model.Session)
|
return wc.session.Load().(*model.Session)
|
||||||
@@ -147,6 +153,12 @@ func (wc *WebConn) Pump() {
|
|||||||
wc.App.HubUnregister(wc)
|
wc.App.HubUnregister(wc)
|
||||||
close(wc.pumpFinished)
|
close(wc.pumpFinished)
|
||||||
|
|
||||||
|
// TODO:
|
||||||
|
// Check if the channel is closed or not,
|
||||||
|
// if closed, then remove the entry from conn manager
|
||||||
|
// else
|
||||||
|
// take both channels, and store them in connection manager.
|
||||||
|
|
||||||
defer ReturnSessionToPool(wc.GetSession())
|
defer ReturnSessionToPool(wc.GetSession())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -195,8 +207,7 @@ func (wc *WebConn) writePump() {
|
|||||||
select {
|
select {
|
||||||
case msg, ok := <-wc.send:
|
case msg, ok := <-wc.send:
|
||||||
if !ok {
|
if !ok {
|
||||||
wc.WebSocket.SetWriteDeadline(time.Now().Add(writeWaitTime))
|
wc.writeMessage(websocket.CloseMessage, []byte{})
|
||||||
wc.WebSocket.WriteMessage(websocket.CloseMessage, []byte{})
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -250,18 +261,18 @@ func (wc *WebConn) writePump() {
|
|||||||
mlog.Warn("websocket.full", logData...)
|
mlog.Warn("websocket.full", logData...)
|
||||||
}
|
}
|
||||||
|
|
||||||
wc.WebSocket.SetWriteDeadline(time.Now().Add(writeWaitTime))
|
if err := wc.writeMessage(websocket.TextMessage, buf.Bytes()); err != nil {
|
||||||
if err := wc.WebSocket.WriteMessage(websocket.TextMessage, buf.Bytes()); err != nil {
|
|
||||||
wc.logSocketErr("websocket.send", err)
|
wc.logSocketErr("websocket.send", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TODO: move to dead queue
|
||||||
|
|
||||||
if wc.App.Metrics() != nil {
|
if wc.App.Metrics() != nil {
|
||||||
wc.App.Metrics().IncrementWebSocketBroadcast(msg.EventType())
|
wc.App.Metrics().IncrementWebSocketBroadcast(msg.EventType())
|
||||||
}
|
}
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
wc.WebSocket.SetWriteDeadline(time.Now().Add(writeWaitTime))
|
if err := wc.writeMessage(websocket.PingMessage, []byte{}); err != nil {
|
||||||
if err := wc.WebSocket.WriteMessage(websocket.PingMessage, []byte{}); err != nil {
|
|
||||||
wc.logSocketErr("websocket.ticker", err)
|
wc.logSocketErr("websocket.ticker", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -279,6 +290,13 @@ func (wc *WebConn) writePump() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// writeMessage is a helper utility that wraps the write to the socket
|
||||||
|
// along with setting the write deadline.
|
||||||
|
func (wc *WebConn) writeMessage(msgType int, data []byte) error {
|
||||||
|
wc.WebSocket.SetWriteDeadline(time.Now().Add(writeWaitTime))
|
||||||
|
return wc.WebSocket.WriteMessage(msgType, data)
|
||||||
|
}
|
||||||
|
|
||||||
// InvalidateCache resets all internal data of the WebConn.
|
// InvalidateCache resets all internal data of the WebConn.
|
||||||
func (wc *WebConn) InvalidateCache() {
|
func (wc *WebConn) InvalidateCache() {
|
||||||
wc.allChannelMembers = nil
|
wc.allChannelMembers = nil
|
||||||
@@ -318,7 +336,11 @@ func (wc *WebConn) IsAuthenticated() bool {
|
|||||||
|
|
||||||
func (wc *WebConn) createHelloMessage() *model.WebSocketEvent {
|
func (wc *WebConn) createHelloMessage() *model.WebSocketEvent {
|
||||||
msg := model.NewWebSocketEvent(model.WEBSOCKET_EVENT_HELLO, "", "", wc.UserId, nil)
|
msg := model.NewWebSocketEvent(model.WEBSOCKET_EVENT_HELLO, "", "", wc.UserId, nil)
|
||||||
msg.Add("server_version", fmt.Sprintf("%v.%v.%v.%v", model.CurrentVersion, model.BuildNumber, wc.App.ClientConfigHash(), wc.App.Srv().License() != nil))
|
msg.Add("server_version", fmt.Sprintf("%v.%v.%v.%v", model.CurrentVersion,
|
||||||
|
model.BuildNumber,
|
||||||
|
wc.App.ClientConfigHash(),
|
||||||
|
wc.App.Srv().License() != nil))
|
||||||
|
msg.Add("connection_id", wc.connectionID.Load())
|
||||||
return msg
|
return msg
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -55,11 +55,12 @@ func (wr *WebSocketRouter) ServeWebSocket(conn *WebConn, r *model.WebSocketReque
|
|||||||
conn.WebSocket.Close()
|
conn.WebSocket.Close()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
conn.SetSession(session)
|
conn.SetSession(session)
|
||||||
conn.SetSessionToken(session.Token)
|
conn.SetSessionToken(session.Token)
|
||||||
conn.UserId = session.UserId
|
conn.UserId = session.UserId
|
||||||
|
|
||||||
|
// TODO: Same logic to reconnect queue as api4/websocket.go
|
||||||
|
|
||||||
wr.app.HubRegister(conn)
|
wr.app.HubRegister(conn)
|
||||||
|
|
||||||
wr.app.Srv().Go(func() {
|
wr.app.Srv().Go(func() {
|
||||||
|
|||||||
@@ -373,6 +373,7 @@ type ServiceSettings struct {
|
|||||||
CollapsedThreads *string `access:"experimental"`
|
CollapsedThreads *string `access:"experimental"`
|
||||||
ManagedResourcePaths *string `access:"environment,write_restrictable,cloud_restrictable"`
|
ManagedResourcePaths *string `access:"environment,write_restrictable,cloud_restrictable"`
|
||||||
EnableLegacySidebar *bool `access:"experimental"`
|
EnableLegacySidebar *bool `access:"experimental"`
|
||||||
|
EnableReliableWebSockets *bool `access:"experimental"` // telemetry: none
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *ServiceSettings) SetDefaults(isUpdate bool) {
|
func (s *ServiceSettings) SetDefaults(isUpdate bool) {
|
||||||
@@ -818,6 +819,10 @@ func (s *ServiceSettings) SetDefaults(isUpdate bool) {
|
|||||||
if s.EnableLegacySidebar == nil {
|
if s.EnableLegacySidebar == nil {
|
||||||
s.EnableLegacySidebar = NewBool(false)
|
s.EnableLegacySidebar = NewBool(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if s.EnableReliableWebSockets == nil {
|
||||||
|
s.EnableReliableWebSockets = NewBool(false)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
type ClusterSettings struct {
|
type ClusterSettings struct {
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user