From 129f0aabd3ea7310f46c3c12e2fa62b029173d67 Mon Sep 17 00:00:00 2001 From: Agniva De Sarker Date: Tue, 7 Dec 2021 21:53:24 +0530 Subject: [PATCH] Remove the remaining fields from *app.Server (#19113) We move these fields to Channels: ``` uploadLockMapMut sync.Mutex uploadLockMap map[string]bool imgDecoder *imaging.Decoder imgEncoder *imaging.Encoder dndTaskMut sync.Mutex dndTask *model.ScheduledTask ``` I think this PR should conclue the initial phase of migrating stuff from Server to Channels. The remaining task would be to focus on continue to create the remaining services from the common things like users, teams, push notifications, clustering for other products to consume. https://community-daily.mattermost.com/boards/workspace/zyoahc9uapdn3xdptac6jb69ic/285b80a3-257d-41f6-8cf4-ed80ca9d92e5/495cdb4d-c13a-4992-8eb9-80cfee2819a4/87df1e15-588e-49ff-8bd1-ffa9651b8c82 ```release-note NONE ``` --- app/brand.go | 4 ++-- app/channels.go | 40 +++++++++++++++++++++++++++++++++++++-- app/file.go | 12 ++++++------ app/server.go | 50 ++++++++----------------------------------------- app/slack.go | 2 +- app/team.go | 2 +- app/upload.go | 16 ++++++++-------- app/user.go | 4 ++-- 8 files changed, 66 insertions(+), 64 deletions(-) diff --git a/app/brand.go b/app/brand.go index 0465dff32a..a5ab57c9d0 100644 --- a/app/brand.go +++ b/app/brand.go @@ -32,13 +32,13 @@ func (a *App) SaveBrandImage(imageData *multipart.FileHeader) *model.AppError { return model.NewAppError("SaveBrandImage", "brand.save_brand_image.check_image_limits.app_error", nil, err.Error(), http.StatusBadRequest) } - img, _, err := a.ch.srv.imgDecoder.Decode(file) + img, _, err := a.ch.imgDecoder.Decode(file) if err != nil { return model.NewAppError("SaveBrandImage", "brand.save_brand_image.decode.app_error", nil, err.Error(), http.StatusBadRequest) } buf := new(bytes.Buffer) - err = a.ch.srv.imgEncoder.EncodePNG(buf, img) + err = a.ch.imgEncoder.EncodePNG(buf, img) if err != nil { return model.NewAppError("SaveBrandImage", "brand.save_brand_image.encode.app_error", nil, err.Error(), http.StatusInternalServerError) } diff --git a/app/channels.go b/app/channels.go index 7005cf0509..03b36fceca 100644 --- a/app/channels.go +++ b/app/channels.go @@ -4,9 +4,11 @@ package app import ( + "runtime" "sync" "sync/atomic" + "github.com/mattermost/mattermost-server/v6/app/imaging" "github.com/mattermost/mattermost-server/v6/app/request" "github.com/mattermost/mattermost-server/v6/einterfaces" "github.com/mattermost/mattermost-server/v6/model" @@ -45,6 +47,18 @@ type Channels struct { Compliance einterfaces.ComplianceInterface DataRetention einterfaces.DataRetentionInterface MessageExport einterfaces.MessageExportInterface + + // These are used to prevent concurrent upload requests + // for a given upload session which could cause inconsistencies + // and data corruption. + uploadLockMapMut sync.Mutex + uploadLockMap map[string]bool + + imgDecoder *imaging.Decoder + imgEncoder *imaging.Encoder + + dndTaskMut sync.Mutex + dndTask *model.ScheduledTask } func init() { @@ -55,8 +69,9 @@ func init() { func NewChannels(s *Server) (*Channels, error) { ch := &Channels{ - srv: s, - imageProxy: imageproxy.MakeImageProxy(s, s.httpService, s.Log), + srv: s, + imageProxy: imageproxy.MakeImageProxy(s, s.httpService, s.Log), + uploadLockMap: map[string]bool{}, } // We are passing a partially filled Channels struct so that the enterprise // methods can have access to app methods. @@ -75,6 +90,20 @@ func NewChannels(s *Server) (*Channels, error) { ch.AccountMigration = accountMigrationInterface(New(ServerConnector(ch))) } + var imgErr error + ch.imgDecoder, imgErr = imaging.NewDecoder(imaging.DecoderOptions{ + ConcurrencyLevel: runtime.NumCPU(), + }) + if imgErr != nil { + return nil, errors.Wrap(imgErr, "failed to create image decoder") + } + ch.imgEncoder, imgErr = imaging.NewEncoder(imaging.EncoderOptions{ + ConcurrencyLevel: runtime.NumCPU(), + }) + if imgErr != nil { + return nil, errors.Wrap(imgErr, "failed to create image encoder") + } + // Setup routes. pluginsRoute := ch.srv.Router.PathPrefix("/plugins/{plugin_id:[A-Za-z0-9\\_\\-\\.]+}").Subrouter() pluginsRoute.HandleFunc("", ch.ServePluginRequest) @@ -109,6 +138,13 @@ func (ch *Channels) Start() error { func (ch *Channels) Stop() error { ch.ShutDownPlugins() + + ch.dndTaskMut.Lock() + if ch.dndTask != nil { + ch.dndTask.Cancel() + } + ch.dndTaskMut.Unlock() + return nil } diff --git a/app/file.go b/app/file.go index 2643424452..ed271ff9be 100644 --- a/app/file.go +++ b/app/file.go @@ -703,8 +703,8 @@ func (a *App) UploadFileX(c *request.Context, channelID, name string, input io.R Input: input, maxFileSize: *a.Config().FileSettings.MaxFileSize, maxImageRes: *a.Config().FileSettings.MaxImageResolution, - imgDecoder: a.ch.srv.imgDecoder, - imgEncoder: a.ch.srv.imgEncoder, + imgDecoder: a.ch.imgDecoder, + imgEncoder: a.ch.imgEncoder, } for _, o := range opts { o(t) @@ -1040,7 +1040,7 @@ func (a *App) HandleImages(previewPathList []string, thumbnailPathList []string, wg := new(sync.WaitGroup) for i := range fileData { - img, release, err := prepareImage(a.ch.srv.imgDecoder, bytes.NewReader(fileData[i])) + img, release, err := prepareImage(a.ch.imgDecoder, bytes.NewReader(fileData[i])) if err != nil { mlog.Debug("Failed to prepare image", mlog.Err(err)) continue @@ -1088,7 +1088,7 @@ func prepareImage(imgDecoder *imaging.Decoder, imgData io.ReadSeeker) (img image func (a *App) generateThumbnailImage(img image.Image, thumbnailPath string) { var buf bytes.Buffer - if err := a.ch.srv.imgEncoder.EncodeJPEG(&buf, imaging.GenerateThumbnail(img, imageThumbnailWidth, imageThumbnailHeight), jpegEncQuality); err != nil { + if err := a.ch.imgEncoder.EncodeJPEG(&buf, imaging.GenerateThumbnail(img, imageThumbnailWidth, imageThumbnailHeight), jpegEncQuality); err != nil { mlog.Error("Unable to encode image as jpeg", mlog.String("path", thumbnailPath), mlog.Err(err)) return } @@ -1103,7 +1103,7 @@ func (a *App) generatePreviewImage(img image.Image, previewPath string) { var buf bytes.Buffer preview := imaging.GeneratePreview(img, imagePreviewWidth) - if err := a.ch.srv.imgEncoder.EncodeJPEG(&buf, preview, jpegEncQuality); err != nil { + if err := a.ch.imgEncoder.EncodeJPEG(&buf, preview, jpegEncQuality); err != nil { mlog.Error("Unable to encode image as preview jpg", mlog.Err(err), mlog.String("path", previewPath)) return } @@ -1124,7 +1124,7 @@ func (a *App) generateMiniPreview(fi *model.FileInfo) { return } defer file.Close() - img, release, err := prepareImage(a.ch.srv.imgDecoder, file) + img, release, err := prepareImage(a.ch.imgDecoder, file) if err != nil { mlog.Debug("generateMiniPreview: prepareImage failed", mlog.Err(err), mlog.String("fileinfo_id", fi.Id), mlog.String("channel_id", fi.ChannelId), diff --git a/app/server.go b/app/server.go index 71d9a3863e..4b1d20d9d9 100644 --- a/app/server.go +++ b/app/server.go @@ -38,7 +38,6 @@ import ( "github.com/mattermost/mattermost-server/v6/app/email" "github.com/mattermost/mattermost-server/v6/app/featureflag" - "github.com/mattermost/mattermost-server/v6/app/imaging" "github.com/mattermost/mattermost-server/v6/app/request" "github.com/mattermost/mattermost-server/v6/app/teams" "github.com/mattermost/mattermost-server/v6/app/users" @@ -174,23 +173,11 @@ type Server struct { tracer *tracing.Tracer - // These are used to prevent concurrent upload requests - // for a given upload session which could cause inconsistencies - // and data corruption. - uploadLockMapMut sync.Mutex - uploadLockMap map[string]bool - featureFlagSynchronizer *featureflag.Synchronizer featureFlagStop chan struct{} featureFlagStopped chan struct{} featureFlagSynchronizerMutex sync.Mutex - imgDecoder *imaging.Decoder - imgEncoder *imaging.Encoder - - dndTaskMut sync.Mutex - dndTask *model.ScheduledTask - products map[string]Product } @@ -207,7 +194,6 @@ func NewServer(options ...Option) (*Server, error) { }, licenseListeners: map[string]func(*model.License, *model.License){}, hashSeed: maphash.MakeSeed(), - uploadLockMap: map[string]bool{}, timezones: timezones.New(), products: make(map[string]Product), } @@ -492,20 +478,6 @@ func NewServer(options ...Option) (*Server, error) { s.setupFeatureFlags() - var imgErr error - s.imgDecoder, imgErr = imaging.NewDecoder(imaging.DecoderOptions{ - ConcurrencyLevel: runtime.NumCPU(), - }) - if imgErr != nil { - return nil, errors.Wrap(imgErr, "failed to create image decoder") - } - s.imgEncoder, imgErr = imaging.NewEncoder(imaging.EncoderOptions{ - ConcurrencyLevel: runtime.NumCPU(), - }) - if imgErr != nil { - return nil, errors.Wrap(imgErr, "failed to create image encoder") - } - s.initJobs() s.clusterLeaderListenerId = s.AddClusterLeaderChangedListener(func() { @@ -1027,12 +999,6 @@ func (s *Server) Shutdown() { } } - s.dndTaskMut.Lock() - if s.dndTask != nil { - s.dndTask.Cancel() - } - s.dndTaskMut.Unlock() - mlog.Info("Server stopped") // Stop products. @@ -2257,18 +2223,18 @@ func (s *Server) ReadFile(path string) ([]byte, *model.AppError) { } func createDNDStatusExpirationRecurringTask(a *App) { - a.ch.srv.dndTaskMut.Lock() - a.ch.srv.dndTask = model.CreateRecurringTaskFromNextIntervalTime("Unset DND Statuses", a.UpdateDNDStatusOfUsers, 5*time.Minute) - a.ch.srv.dndTaskMut.Unlock() + a.ch.dndTaskMut.Lock() + a.ch.dndTask = model.CreateRecurringTaskFromNextIntervalTime("Unset DND Statuses", a.UpdateDNDStatusOfUsers, 5*time.Minute) + a.ch.dndTaskMut.Unlock() } func cancelDNDStatusExpirationRecurringTask(a *App) { - a.ch.srv.dndTaskMut.Lock() - if a.ch.srv.dndTask != nil { - a.ch.srv.dndTask.Cancel() - a.ch.srv.dndTask = nil + a.ch.dndTaskMut.Lock() + if a.ch.dndTask != nil { + a.ch.dndTask.Cancel() + a.ch.dndTask = nil } - a.ch.srv.dndTaskMut.Unlock() + a.ch.dndTaskMut.Unlock() } func runDNDStatusExpireJob(a *App) { diff --git a/app/slack.go b/app/slack.go index 58c1e2f32d..b355396955 100644 --- a/app/slack.go +++ b/app/slack.go @@ -41,7 +41,7 @@ func (a *App) SlackImport(c *request.Context, fileData multipart.File, fileSize InvalidateAllCaches: func() { a.ch.srv.InvalidateAllCaches() }, MaxPostSize: func() int { return a.ch.srv.MaxPostSize() }, PrepareImage: func(fileData []byte) (image.Image, func(), error) { - img, release, err := prepareImage(a.ch.srv.imgDecoder, bytes.NewReader(fileData)) + img, release, err := prepareImage(a.ch.imgDecoder, bytes.NewReader(fileData)) if err != nil { return nil, nil, err } diff --git a/app/team.go b/app/team.go index 3eed761f2a..eb8acb12ad 100644 --- a/app/team.go +++ b/app/team.go @@ -1848,7 +1848,7 @@ func (a *App) SetTeamIconFromFile(team *model.Team, file io.Reader) *model.AppEr img = imaging.FillCenter(img, teamIconWidthAndHeight, teamIconWidthAndHeight) buf := new(bytes.Buffer) - err = a.Srv().imgEncoder.EncodePNG(buf, img) + err = a.ch.imgEncoder.EncodePNG(buf, img) if err != nil { return model.NewAppError("SetTeamIcon", "api.team.set_team_icon.encode.app_error", nil, err.Error(), http.StatusInternalServerError) } diff --git a/app/upload.go b/app/upload.go index 5470a64687..d472a6755a 100644 --- a/app/upload.go +++ b/app/upload.go @@ -166,23 +166,23 @@ func (a *App) GetUploadSessionsForUser(userID string) ([]*model.UploadSession, * func (a *App) UploadData(c *request.Context, us *model.UploadSession, rd io.Reader) (*model.FileInfo, *model.AppError) { // prevent more than one caller to upload data at the same time for a given upload session. // This is to avoid possible inconsistencies. - a.Srv().uploadLockMapMut.Lock() - locked := a.Srv().uploadLockMap[us.Id] + a.ch.uploadLockMapMut.Lock() + locked := a.ch.uploadLockMap[us.Id] if locked { // session lock is already taken, return error. - a.Srv().uploadLockMapMut.Unlock() + a.ch.uploadLockMapMut.Unlock() return nil, model.NewAppError("UploadData", "app.upload.upload_data.concurrent.app_error", nil, "", http.StatusBadRequest) } // grab the session lock. - a.Srv().uploadLockMap[us.Id] = true - a.Srv().uploadLockMapMut.Unlock() + a.ch.uploadLockMap[us.Id] = true + a.ch.uploadLockMapMut.Unlock() // reset the session lock on exit. defer func() { - a.Srv().uploadLockMapMut.Lock() - delete(a.Srv().uploadLockMap, us.Id) - a.Srv().uploadLockMapMut.Unlock() + a.ch.uploadLockMapMut.Lock() + delete(a.ch.uploadLockMap, us.Id) + a.ch.uploadLockMapMut.Unlock() }() // fetch the session from store to check for inconsistencies. diff --git a/app/user.go b/app/user.go index 54e2f69566..7ed9fc783a 100644 --- a/app/user.go +++ b/app/user.go @@ -767,7 +767,7 @@ func (a *App) SetProfileImageFromMultiPartFile(userID string, file multipart.Fil func (a *App) AdjustImage(file io.Reader) (*bytes.Buffer, *model.AppError) { // Decode image into Image object - img, _, err := a.ch.srv.imgDecoder.Decode(file) + img, _, err := a.ch.imgDecoder.Decode(file) if err != nil { return nil, model.NewAppError("SetProfileImage", "api.user.upload_profile_user.decode.app_error", nil, err.Error(), http.StatusBadRequest) } @@ -780,7 +780,7 @@ func (a *App) AdjustImage(file io.Reader) (*bytes.Buffer, *model.AppError) { img = imaging.FillCenter(img, profileWidthAndHeight, profileWidthAndHeight) buf := new(bytes.Buffer) - err = a.ch.srv.imgEncoder.EncodePNG(buf, img) + err = a.ch.imgEncoder.EncodePNG(buf, img) if err != nil { return nil, model.NewAppError("SetProfileImage", "api.user.upload_profile_user.encode.app_error", nil, err.Error(), http.StatusInternalServerError) }