Optimizing weh_hub.Start() for master (#5576)
* Optimizing weh_hub.Start() * Optimizing weh_hub.Start() * Optimizing weh_hub.Start() * Adding to IsAuthenticated * Fixing problem with IsAuthenticated
Этот коммит содержится в:
@@ -29,6 +29,7 @@ type WebConn struct {
|
|||||||
Send chan model.WebSocketMessage
|
Send chan model.WebSocketMessage
|
||||||
SessionToken string
|
SessionToken string
|
||||||
SessionExpiresAt int64
|
SessionExpiresAt int64
|
||||||
|
Session *model.Session
|
||||||
UserId string
|
UserId string
|
||||||
T goi18n.TranslateFunc
|
T goi18n.TranslateFunc
|
||||||
Locale string
|
Locale string
|
||||||
@@ -148,6 +149,7 @@ func (webCon *WebConn) InvalidateCache() {
|
|||||||
webCon.AllChannelMembers = nil
|
webCon.AllChannelMembers = nil
|
||||||
webCon.LastAllChannelMembersTime = 0
|
webCon.LastAllChannelMembersTime = 0
|
||||||
webCon.SessionExpiresAt = 0
|
webCon.SessionExpiresAt = 0
|
||||||
|
webCon.Session = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (webCon *WebConn) IsAuthenticated() bool {
|
func (webCon *WebConn) IsAuthenticated() bool {
|
||||||
@@ -162,11 +164,13 @@ func (webCon *WebConn) IsAuthenticated() bool {
|
|||||||
l4g.Error(utils.T("api.websocket.invalid_session.error"), err.Error())
|
l4g.Error(utils.T("api.websocket.invalid_session.error"), err.Error())
|
||||||
webCon.SessionToken = ""
|
webCon.SessionToken = ""
|
||||||
webCon.SessionExpiresAt = 0
|
webCon.SessionExpiresAt = 0
|
||||||
|
webCon.Session = nil
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
webCon.SessionToken = session.Token
|
webCon.SessionToken = session.Token
|
||||||
webCon.SessionExpiresAt = session.ExpiresAt
|
webCon.SessionExpiresAt = session.ExpiresAt
|
||||||
|
webCon.Session = session
|
||||||
}
|
}
|
||||||
|
|
||||||
return true
|
return true
|
||||||
@@ -231,17 +235,23 @@ func (webCon *WebConn) ShouldSendEvent(msg *model.WebSocketEvent) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (webCon *WebConn) IsMemberOfTeam(teamId string) bool {
|
func (webCon *WebConn) IsMemberOfTeam(teamId string) bool {
|
||||||
session, err := GetSession(webCon.SessionToken)
|
|
||||||
if err != nil {
|
|
||||||
l4g.Error(utils.T("api.websocket.invalid_session.error"), err.Error())
|
|
||||||
return false
|
|
||||||
} else {
|
|
||||||
member := session.GetTeamByTeamId(teamId)
|
|
||||||
|
|
||||||
if member != nil {
|
if webCon.Session == nil {
|
||||||
return true
|
session, err := GetSession(webCon.SessionToken)
|
||||||
} else {
|
if err != nil {
|
||||||
|
l4g.Error(utils.T("api.websocket.invalid_session.error"), err.Error())
|
||||||
return false
|
return false
|
||||||
|
} else {
|
||||||
|
webCon.Session = session
|
||||||
}
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
member := webCon.Session.GetTeamByTeamId(teamId)
|
||||||
|
|
||||||
|
if member != nil {
|
||||||
|
return true
|
||||||
|
} else {
|
||||||
|
return false
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -89,6 +89,15 @@ func HubUnregister(webConn *WebConn) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func Publish(message *model.WebSocketEvent) {
|
func Publish(message *model.WebSocketEvent) {
|
||||||
|
|
||||||
|
if SkipTypingMessage(message) {
|
||||||
|
if metrics := einterfaces.GetMetricsInterface(); metrics != nil {
|
||||||
|
metrics.IncrementWebsocketEvent(message.Event + "_skipped")
|
||||||
|
}
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
if metrics := einterfaces.GetMetricsInterface(); metrics != nil {
|
if metrics := einterfaces.GetMetricsInterface(); metrics != nil {
|
||||||
metrics.IncrementWebsocketEvent(message.Event)
|
metrics.IncrementWebsocketEvent(message.Event)
|
||||||
}
|
}
|
||||||
@@ -278,20 +287,18 @@ func (h *Hub) Start() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
case msg := <-h.broadcast:
|
case msg := <-h.broadcast:
|
||||||
if OkToSendTypingMessage(msg) {
|
for _, webCon := range h.connections {
|
||||||
for _, webCon := range h.connections {
|
if webCon.ShouldSendEvent(msg) {
|
||||||
if webCon.ShouldSendEvent(msg) {
|
select {
|
||||||
select {
|
case webCon.Send <- msg:
|
||||||
case webCon.Send <- msg:
|
default:
|
||||||
default:
|
l4g.Error(fmt.Sprintf("webhub.broadcast: cannot send, closing websocket for userId=%v", webCon.UserId))
|
||||||
l4g.Error(fmt.Sprintf("webhub.broadcast: cannot send, closing websocket for userId=%v", webCon.UserId))
|
close(webCon.Send)
|
||||||
close(webCon.Send)
|
for i, webConCandidate := range h.connections {
|
||||||
for i, webConCandidate := range h.connections {
|
if webConCandidate == webCon {
|
||||||
if webConCandidate == webCon {
|
h.connections[i] = h.connections[len(h.connections)-1]
|
||||||
h.connections[i] = h.connections[len(h.connections)-1]
|
h.connections = h.connections[:len(h.connections)-1]
|
||||||
h.connections = h.connections[:len(h.connections)-1]
|
break
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -331,13 +338,13 @@ func (h *Hub) Start() {
|
|||||||
go doRecoverableStart()
|
go doRecoverableStart()
|
||||||
}
|
}
|
||||||
|
|
||||||
func OkToSendTypingMessage(msg *model.WebSocketEvent) bool {
|
func SkipTypingMessage(msg *model.WebSocketEvent) bool {
|
||||||
// Only broadcast typing messages if less than 1K people in channel
|
// Only broadcast typing messages if less than 1K people in channel
|
||||||
if msg.Event == model.WEBSOCKET_EVENT_TYPING {
|
if msg.Event == model.WEBSOCKET_EVENT_TYPING {
|
||||||
if Srv.Store.Channel().GetMemberCountFromCache(msg.Broadcast.ChannelId) > *utils.Cfg.TeamSettings.MaxNotificationsPerChannel {
|
if Srv.Store.Channel().GetMemberCountFromCache(msg.Broadcast.ChannelId) > *utils.Cfg.TeamSettings.MaxNotificationsPerChannel {
|
||||||
return false
|
return true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return true
|
return false
|
||||||
}
|
}
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user