From e13d85d8c7fa98e1f25052e6d3de8d94c16853b7 Mon Sep 17 00:00:00 2001 From: Ibrahim Serdar Acikgoz Date: Thu, 1 Apr 2021 11:29:56 +0300 Subject: [PATCH] [MM-34179] app: move cluster registeration to server creation (#17306) * app: move cluster registeration to server creation * initserver via fakeapp * reflect review comments --- app/admin.go | 6 +- app/app.go | 2 +- app/app_iface.go | 4 - app/cluster_handlers.go | 144 ++++++++++++++++++--------- app/opentracing/opentracing_layer.go | 45 --------- app/server.go | 1 + app/session.go | 19 +--- app/web_hub.go | 36 +------ cmd/mattermost/commands/server.go | 7 ++ 9 files changed, 116 insertions(+), 148 deletions(-) diff --git a/app/admin.go b/app/admin.go index 7b4c9d30e2..de743d1489 100644 --- a/app/admin.go +++ b/app/admin.go @@ -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 { diff --git a/app/app.go b/app/app.go index d868815967..b69f861884 100644 --- a/app/app.go +++ b/app/app.go @@ -76,7 +76,7 @@ func (a *App) InitServer() { a.initJobs() if a.srv.joinCluster && a.srv.Cluster != nil { - a.registerAllClusterMessageHandlers() + a.registerAppClusterMessageHandlers() } a.DoAppMigrations() diff --git a/app/app_iface.go b/app/app_iface.go index 947482e9d8..5e00409720 100644 --- a/app/app_iface.go +++ b/app/app_iface.go @@ -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 diff --git a/app/cluster_handlers.go b/app/cluster_handlers.go index 1d3a7578a2..9adf95b4e4 100644 --- a/app/cluster_handlers.go +++ b/app/cluster_handlers.go @@ -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) + } } diff --git a/app/opentracing/opentracing_layer.go b/app/opentracing/opentracing_layer.go index de7016920c..0179055d05 100644 --- a/app/opentracing/opentracing_layer.go +++ b/app/opentracing/opentracing_layer.go @@ -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") diff --git a/app/server.go b/app/server.go index 5ec7e23191..95b3638b84 100644 --- a/app/server.go +++ b/app/server.go @@ -454,6 +454,7 @@ func NewServer(options ...Option) (*Server, error) { }) if s.joinCluster && s.Cluster != nil { + s.registerClusterHandlers() s.Cluster.StartInterNodeCommunication() } diff --git a/app/session.go b/app/session.go index 33b97b67cd..f2f8da2a71 100644 --- a/app/session.go +++ b/app/session.go @@ -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) { diff --git a/app/web_hub.go b/app/web_hub.go index 209109a3da..fbdddc3e62 100644 --- a/app/web_hub.go +++ b/app/web_hub.go @@ -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 { diff --git a/cmd/mattermost/commands/server.go b/cmd/mattermost/commands/server.go index fb102a03e5..ae87586b36 100644 --- a/cmd/mattermost/commands/server.go +++ b/cmd/mattermost/commands/server.go @@ -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)