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) }