[MM-34179] app: move cluster registeration to server creation (#17306)

* app: move cluster registeration to server creation

* initserver via fakeapp

* reflect review comments
Этот коммит содержится в:
Ibrahim Serdar Acikgoz
2021-04-01 11:29:56 +03:00
коммит произвёл GitHub
родитель 756f5fbff3
Коммит e13d85d8c7
9 изменённых файлов: 116 добавлений и 148 удалений

Просмотреть файл

@@ -235,9 +235,9 @@ func (a *App) TestEmail(userID string, cfg *model.Config) *model.AppError {
return nil
}
// ServerBusyStateChanged is called when a CLUSTER_EVENT_BUSY_STATE_CHANGED is received.
func (a *App) ServerBusyStateChanged(sbs *model.ServerBusyState) {
a.Srv().Busy.ClusterEventChanged(sbs)
// serverBusyStateChanged is called when a CLUSTER_EVENT_BUSY_STATE_CHANGED is received.
func (s *Server) serverBusyStateChanged(sbs *model.ServerBusyState) {
s.Busy.ClusterEventChanged(sbs)
if sbs.Busy {
mlog.Warn("server busy state activitated via cluster event - non-critical services disabled", mlog.Int64("expires_sec", sbs.Expires))
} else {

Просмотреть файл

@@ -76,7 +76,7 @@ func (a *App) InitServer() {
a.initJobs()
if a.srv.joinCluster && a.srv.Cluster != nil {
a.registerAllClusterMessageHandlers()
a.registerAppClusterMessageHandlers()
}
a.DoAppMigrations()

Просмотреть файл

@@ -276,8 +276,6 @@ type AppIface interface {
// ServePluginPublicRequest serves public plugin files
// at the URL http(s)://$SITE_URL/plugins/$PLUGIN_ID/public/{anything}
ServePluginPublicRequest(w http.ResponseWriter, r *http.Request)
// ServerBusyStateChanged is called when a CLUSTER_EVENT_BUSY_STATE_CHANGED is received.
ServerBusyStateChanged(sbs *model.ServerBusyState)
// SessionHasPermissionToManageBot returns nil if the session has access to manage the given bot.
// This function deviates from other authorization checks in returning an error instead of just
// a boolean, allowing the permission failure to be exposed with more granularity.
@@ -768,7 +766,6 @@ type AppIface interface {
InstallPluginFromData(data model.PluginEventData)
InvalidateAllEmailInvites() *model.AppError
InvalidateCacheForUser(userID string)
InvalidateWebConnSessionCacheForUser(userID string)
InviteGuestsToChannels(teamID string, guestsInvite *model.GuestsInvite, senderId string) *model.AppError
InviteGuestsToChannelsGracefully(teamID string, guestsInvite *model.GuestsInvite, senderId string) ([]*model.EmailInviteWithError, *model.AppError)
InviteNewUsersToTeam(emailList []string, teamID, senderId string) *model.AppError
@@ -836,7 +833,6 @@ type AppIface interface {
PreparePostListForClient(originalList *model.PostList) *model.PostList
ProcessSlackText(text string) string
Publish(message *model.WebSocketEvent)
PublishSkipClusterSend(message *model.WebSocketEvent)
PublishUserTyping(userID, channelID, parentId string) *model.AppError
PurgeBleveIndexes() *model.AppError
PurgeElasticsearchIndexes() *model.AppError

Просмотреть файл

@@ -7,59 +7,19 @@ import (
"strings"
"github.com/mattermost/mattermost-server/v5/model"
"github.com/mattermost/mattermost-server/v5/shared/mlog"
)
// RegisterAllClusterMessageHandlers registers the cluster message handlers that are handled by the App layer.
// registerAppClusterMessageHandlers registers the cluster message handlers that are handled by the App layer.
//
// The cluster event handlers are spread across this function and
// The cluster event handlers are spread across this function, Server.registerClusterHandlers and
// NewLocalCacheLayer. Be careful to not have duplicated handlers here and
// there.
func (a *App) registerAllClusterMessageHandlers() {
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_PUBLISH, a.clusterPublishHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_UPDATE_STATUS, a.clusterUpdateStatusHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_ALL_CACHES, a.clusterInvalidateAllCachesHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_CHANNEL_MEMBERS_NOTIFY_PROPS, a.clusterInvalidateCacheForChannelMembersNotifyPropHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_CHANNEL_BY_NAME, a.clusterInvalidateCacheForChannelByNameHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_USER, a.clusterInvalidateCacheForUserHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_USER_TEAMS, a.clusterInvalidateCacheForUserTeamsHandler)
func (a *App) registerAppClusterMessageHandlers() {
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_CLEAR_SESSION_CACHE_FOR_USER, a.clusterClearSessionCacheForUserHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_CLEAR_SESSION_CACHE_FOR_ALL_USERS, a.clusterClearSessionCacheForAllUsersHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_INSTALL_PLUGIN, a.clusterInstallPluginHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_REMOVE_PLUGIN, a.clusterRemovePluginHandler)
a.Cluster().RegisterClusterMessageHandler(model.CLUSTER_EVENT_BUSY_STATE_CHANGED, a.clusterBusyStateChgHandler)
}
func (a *App) clusterPublishHandler(msg *model.ClusterMessage) {
event := model.WebSocketEventFromJson(strings.NewReader(msg.Data))
if event == nil {
return
}
a.PublishSkipClusterSend(event)
}
func (a *App) clusterUpdateStatusHandler(msg *model.ClusterMessage) {
status := model.StatusFromJson(strings.NewReader(msg.Data))
a.AddStatusCacheSkipClusterSend(status)
}
func (a *App) clusterInvalidateAllCachesHandler(msg *model.ClusterMessage) {
a.Srv().InvalidateAllCachesSkipSend()
}
func (a *App) clusterInvalidateCacheForChannelMembersNotifyPropHandler(msg *model.ClusterMessage) {
a.invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(msg.Data)
}
func (a *App) clusterInvalidateCacheForChannelByNameHandler(msg *model.ClusterMessage) {
a.invalidateCacheForChannelByNameSkipClusterSend(msg.Props["id"], msg.Props["name"])
}
func (a *App) clusterInvalidateCacheForUserHandler(msg *model.ClusterMessage) {
a.invalidateCacheForUserSkipClusterSend(msg.Data)
}
func (a *App) clusterInvalidateCacheForUserTeamsHandler(msg *model.ClusterMessage) {
a.InvalidateWebConnSessionCacheForUser(msg.Data)
}
func (a *App) clusterClearSessionCacheForUserHandler(msg *model.ClusterMessage) {
@@ -78,6 +38,98 @@ func (a *App) clusterRemovePluginHandler(msg *model.ClusterMessage) {
a.RemovePluginFromData(model.PluginEventDataFromJson(strings.NewReader(msg.Data)))
}
func (a *App) clusterBusyStateChgHandler(msg *model.ClusterMessage) {
a.ServerBusyStateChanged(model.ServerBusyStateFromJson(strings.NewReader(msg.Data)))
// registerClusterHandlers registers the cluster message handlers that are handled by the server.
func (s *Server) registerClusterHandlers() {
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_PUBLISH, s.clusterPublishHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_UPDATE_STATUS, s.clusterUpdateStatusHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_ALL_CACHES, s.clusterInvalidateAllCachesHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_CHANNEL_MEMBERS_NOTIFY_PROPS, s.clusterInvalidateCacheForChannelMembersNotifyPropHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_CHANNEL_BY_NAME, s.clusterInvalidateCacheForChannelByNameHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_USER, s.clusterInvalidateCacheForUserHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_INVALIDATE_CACHE_FOR_USER_TEAMS, s.clusterInvalidateCacheForUserTeamsHandler)
s.Cluster.RegisterClusterMessageHandler(model.CLUSTER_EVENT_BUSY_STATE_CHANGED, s.clusterBusyStateChgHandler)
}
func (s *Server) clusterPublishHandler(msg *model.ClusterMessage) {
event := model.WebSocketEventFromJson(strings.NewReader(msg.Data))
if event == nil {
return
}
s.PublishSkipClusterSend(event)
}
func (s *Server) clusterUpdateStatusHandler(msg *model.ClusterMessage) {
status := model.StatusFromJson(strings.NewReader(msg.Data))
s.statusCache.Set(status.UserId, status)
}
func (s *Server) clusterInvalidateAllCachesHandler(msg *model.ClusterMessage) {
s.InvalidateAllCachesSkipSend()
}
func (s *Server) clusterInvalidateCacheForChannelMembersNotifyPropHandler(msg *model.ClusterMessage) {
s.invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(msg.Data)
}
func (s *Server) clusterInvalidateCacheForChannelByNameHandler(msg *model.ClusterMessage) {
s.invalidateCacheForChannelByNameSkipClusterSend(msg.Props["id"], msg.Props["name"])
}
func (s *Server) clusterInvalidateCacheForUserHandler(msg *model.ClusterMessage) {
s.invalidateCacheForUserSkipClusterSend(msg.Data)
}
func (s *Server) clusterInvalidateCacheForUserTeamsHandler(msg *model.ClusterMessage) {
s.invalidateWebConnSessionCacheForUser(msg.Data)
}
func (s *Server) clearSessionCacheForUserSkipClusterSend(userID string) {
if keys, err := s.sessionCache.Keys(); err == nil {
var session *model.Session
for _, key := range keys {
if err := s.sessionCache.Get(key, &session); err == nil {
if session.UserId == userID {
s.sessionCache.Remove(key)
if s.Metrics != nil {
s.Metrics.IncrementMemCacheInvalidationCounterSession()
}
}
}
}
}
s.invalidateWebConnSessionCacheForUser(userID)
}
func (s *Server) clearSessionCacheForAllUsersSkipClusterSend() {
mlog.Info("Purging sessions cache")
s.sessionCache.Purge()
}
func (s *Server) clusterBusyStateChgHandler(msg *model.ClusterMessage) {
s.serverBusyStateChanged(model.ServerBusyStateFromJson(strings.NewReader(msg.Data)))
}
func (s *Server) invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(channelID string) {
s.Store.Channel().InvalidateCacheForChannelMembersNotifyProps(channelID)
}
func (s *Server) invalidateCacheForChannelByNameSkipClusterSend(teamID, name string) {
if teamID == "" {
teamID = "dm"
}
s.Store.Channel().InvalidateChannelByName(teamID, name)
}
func (s *Server) invalidateCacheForUserSkipClusterSend(userID string) {
s.Store.Channel().InvalidateAllChannelMembersForUser(userID)
s.invalidateWebConnSessionCacheForUser(userID)
}
func (s *Server) invalidateWebConnSessionCacheForUser(userID string) {
hub := s.GetHubForUserId(userID)
if hub != nil {
hub.InvalidateUser(userID)
}
}

Просмотреть файл

@@ -10012,21 +10012,6 @@ func (a *OpenTracingAppLayer) InvalidateCacheForUser(userID string) {
a.app.InvalidateCacheForUser(userID)
}
func (a *OpenTracingAppLayer) InvalidateWebConnSessionCacheForUser(userID string) {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.InvalidateWebConnSessionCacheForUser")
a.ctx = newCtx
a.app.Srv().Store.SetContext(newCtx)
defer func() {
a.app.Srv().Store.SetContext(origCtx)
a.ctx = origCtx
}()
defer span.Finish()
a.app.InvalidateWebConnSessionCacheForUser(userID)
}
func (a *OpenTracingAppLayer) InviteGuestsToChannels(teamID string, guestsInvite *model.GuestsInvite, senderId string) *model.AppError {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.InviteGuestsToChannels")
@@ -11609,21 +11594,6 @@ func (a *OpenTracingAppLayer) Publish(message *model.WebSocketEvent) {
a.app.Publish(message)
}
func (a *OpenTracingAppLayer) PublishSkipClusterSend(message *model.WebSocketEvent) {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.PublishSkipClusterSend")
a.ctx = newCtx
a.app.Srv().Store.SetContext(newCtx)
defer func() {
a.app.Srv().Store.SetContext(origCtx)
a.ctx = origCtx
}()
defer span.Finish()
a.app.PublishSkipClusterSend(message)
}
func (a *OpenTracingAppLayer) PublishUserTyping(userID string, channelID string, parentId string) *model.AppError {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.PublishUserTyping")
@@ -13540,21 +13510,6 @@ func (a *OpenTracingAppLayer) ServePluginRequest(w http.ResponseWriter, r *http.
a.app.ServePluginRequest(w, r)
}
func (a *OpenTracingAppLayer) ServerBusyStateChanged(sbs *model.ServerBusyState) {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.ServerBusyStateChanged")
a.ctx = newCtx
a.app.Srv().Store.SetContext(newCtx)
defer func() {
a.app.Srv().Store.SetContext(origCtx)
a.ctx = origCtx
}()
defer span.Finish()
a.app.ServerBusyStateChanged(sbs)
}
func (a *OpenTracingAppLayer) SessionCacheLength() int {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SessionCacheLength")

Просмотреть файл

@@ -454,6 +454,7 @@ func NewServer(options ...Option) (*Server, error) {
})
if s.joinCluster && s.Cluster != nil {
s.registerClusterHandlers()
s.Cluster.StartInterNodeCommunication()
}

Просмотреть файл

@@ -237,26 +237,11 @@ func (a *App) ClearSessionCacheForAllUsers() {
}
func (a *App) ClearSessionCacheForUserSkipClusterSend(userID string) {
if keys, err := a.Srv().sessionCache.Keys(); err == nil {
var session *model.Session
for _, key := range keys {
if err := a.Srv().sessionCache.Get(key, &session); err == nil {
if session.UserId == userID {
a.Srv().sessionCache.Remove(key)
if a.Metrics() != nil {
a.Metrics().IncrementMemCacheInvalidationCounterSession()
}
}
}
}
}
a.InvalidateWebConnSessionCacheForUser(userID)
a.Srv().clearSessionCacheForUserSkipClusterSend(userID)
}
func (a *App) ClearSessionCacheForAllUsersSkipClusterSend() {
mlog.Info("Purging sessions cache")
a.Srv().sessionCache.Purge()
a.Srv().clearSessionCacheForAllUsersSkipClusterSend()
}
func (a *App) AddSessionToCache(session *model.Session) {

Просмотреть файл

@@ -94,22 +94,10 @@ func (a *App) HubStart() {
a.srv.hubs = hubs
}
func (a *App) invalidateCacheForUserSkipClusterSend(userID string) {
a.Srv().Store.Channel().InvalidateAllChannelMembersForUser(userID)
a.InvalidateWebConnSessionCacheForUser(userID)
}
func (a *App) invalidateCacheForWebhook(webhookID string) {
a.Srv().Store.Webhook().InvalidateWebhookCache(webhookID)
}
func (a *App) InvalidateWebConnSessionCacheForUser(userID string) {
hub := a.GetHubForUserId(userID)
if hub != nil {
hub.InvalidateUser(userID)
}
}
// HubStop stops all the hubs.
func (s *Server) HubStop() {
mlog.Info("stopping websocket hub connections")
@@ -205,13 +193,9 @@ func (s *Server) PublishSkipClusterSend(message *model.WebSocketEvent) {
}
}
func (a *App) PublishSkipClusterSend(message *model.WebSocketEvent) {
a.Srv().PublishSkipClusterSend(message)
}
func (a *App) invalidateCacheForChannel(channel *model.Channel) {
a.Srv().Store.Channel().InvalidateChannel(channel.Id)
a.invalidateCacheForChannelByNameSkipClusterSend(channel.TeamId, channel.Name)
a.Srv().invalidateCacheForChannelByNameSkipClusterSend(channel.TeamId, channel.Name)
if a.Cluster() != nil {
nameMsg := &model.ClusterMessage{
@@ -238,7 +222,7 @@ func (a *App) invalidateCacheForChannelMembers(channelID string) {
}
func (a *App) invalidateCacheForChannelMembersNotifyProps(channelID string) {
a.invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(channelID)
a.Srv().invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(channelID)
if a.Cluster() != nil {
msg := &model.ClusterMessage{
@@ -250,25 +234,13 @@ func (a *App) invalidateCacheForChannelMembersNotifyProps(channelID string) {
}
}
func (a *App) invalidateCacheForChannelMembersNotifyPropsSkipClusterSend(channelID string) {
a.Srv().Store.Channel().InvalidateCacheForChannelMembersNotifyProps(channelID)
}
func (a *App) invalidateCacheForChannelByNameSkipClusterSend(teamID, name string) {
if teamID == "" {
teamID = "dm"
}
a.Srv().Store.Channel().InvalidateChannelByName(teamID, name)
}
func (a *App) invalidateCacheForChannelPosts(channelID string) {
a.Srv().Store.Channel().InvalidatePinnedPostCount(channelID)
a.Srv().Store.Post().InvalidateLastPostTimeCache(channelID)
}
func (a *App) InvalidateCacheForUser(userID string) {
a.invalidateCacheForUserSkipClusterSend(userID)
a.Srv().invalidateCacheForUserSkipClusterSend(userID)
a.Srv().Store.User().InvalidateProfilesInChannelCacheByUser(userID)
a.Srv().Store.User().InvalidateProfileCacheForUser(userID)
@@ -284,7 +256,7 @@ func (a *App) InvalidateCacheForUser(userID string) {
}
func (a *App) invalidateCacheForUserTeams(userID string) {
a.InvalidateWebConnSessionCacheForUser(userID)
a.Srv().invalidateWebConnSessionCacheForUser(userID)
a.Srv().Store.Team().InvalidateAllTeamIdsForUser(userID)
if a.Cluster() != nil {

Просмотреть файл

@@ -111,6 +111,13 @@ func runServer(configStore *config.Store, usedPlatform bool, interruptChan chan
return serverErr
}
// TODO: remove this and handle all required initialization while creating
// the server. In theory, we shouldn't depend on App to have a fully-featured
// server. This initialization is added so that cluster handlers are registered
// and job schedulers are initialized.
fakeApp := app.New(app.ServerConnector(server))
fakeApp.InitServer()
// If we allow testing then listen for manual testing URL hits
if *server.Config().ServiceSettings.EnableTesting {
manualtesting.Init(api)