From 200a56fa5a2d93ca0097b6ee279170f98079d7b2 Mon Sep 17 00:00:00 2001 From: Claudio Costa Date: Mon, 25 Jan 2021 10:40:30 +0100 Subject: [PATCH] [MM-28423] Implement ImportDelete job (#16588) * Implement ImportDelete job * Add missing translation * Improve logging * Avoid deleting the file in case of errors --- app/app.go | 3 + app/app_iface.go | 1 + app/enterprise.go | 6 + app/file.go | 13 ++ app/import.go | 2 +- app/opentracing/opentracing_layer.go | 22 +++ app/upload.go | 6 +- i18n/en.json | 4 + imports/placeholder.go | 3 + jobs/import_delete/scheduler.go | 51 ++++++ jobs/import_delete/worker.go | 172 +++++++++++++++++++++ jobs/interfaces/import_delete_interface.go | 11 ++ jobs/jobs_watcher.go | 7 + jobs/schedulers.go | 5 + jobs/server.go | 1 + jobs/workers.go | 13 ++ model/job.go | 2 + services/filesstore/filesstore.go | 2 + services/filesstore/filesstore_test.go | 37 +++++ services/filesstore/localstore.go | 9 ++ services/filesstore/mocks/FileBackend.go | 23 +++ services/filesstore/s3store.go | 12 ++ 22 files changed, 401 insertions(+), 4 deletions(-) create mode 100644 jobs/import_delete/scheduler.go create mode 100644 jobs/import_delete/worker.go create mode 100644 jobs/interfaces/import_delete_interface.go diff --git a/app/app.go b/app/app.go index b38ff1588d..2578286f8b 100644 --- a/app/app.go +++ b/app/app.go @@ -119,6 +119,9 @@ func (a *App) initJobs() { if jobsImportProcessInterface != nil { a.srv.Jobs.ImportProcess = jobsImportProcessInterface(a) } + if jobsImportDeleteInterface != nil { + a.srv.Jobs.ImportDelete = jobsImportDeleteInterface(a) + } if jobsActiveUsersInterface != nil { a.srv.Jobs.ActiveUsers = jobsActiveUsersInterface(a) diff --git a/app/app_iface.go b/app/app_iface.go index 697c17c803..21d15f56c2 100644 --- a/app/app_iface.go +++ b/app/app_iface.go @@ -498,6 +498,7 @@ type AppIface interface { FetchSamlMetadataFromIdp(url string) ([]byte, *model.AppError) FileBackend() (filesstore.FileBackend, *model.AppError) FileExists(path string) (bool, *model.AppError) + FileModTime(path string) (time.Time, *model.AppError) FileSize(path string) (int64, *model.AppError) FillInChannelProps(channel *model.Channel) *model.AppError FillInChannelsProps(channelList *model.ChannelList) *model.AppError diff --git a/app/enterprise.go b/app/enterprise.go index 523d31fc40..44e7638c5e 100644 --- a/app/enterprise.go +++ b/app/enterprise.go @@ -114,6 +114,12 @@ func RegisterJobsImportProcessInterface(f func(*App) tjobs.ImportProcessInterfac jobsImportProcessInterface = f } +var jobsImportDeleteInterface func(*App) tjobs.ImportDeleteInterface + +func RegisterJobsImportDeleteInterface(f func(*App) tjobs.ImportDeleteInterface) { + jobsImportDeleteInterface = f +} + var productNoticesJobInterface func(*App) tjobs.ProductNoticesJobInterface func RegisterProductNoticesJobInterface(f func(*App) tjobs.ProductNoticesJobInterface) { diff --git a/app/file.go b/app/file.go index 84f416730e..7acf234be3 100644 --- a/app/file.go +++ b/app/file.go @@ -158,6 +158,19 @@ func (a *App) FileSize(path string) (int64, *model.AppError) { return size, nil } +func (a *App) FileModTime(path string) (time.Time, *model.AppError) { + backend, err := a.FileBackend() + if err != nil { + return time.Time{}, err + } + modTime, nErr := backend.FileModTime(path) + if nErr != nil { + return time.Time{}, model.NewAppError("FileModTime", "api.file.file_mod_time.app_error", nil, nErr.Error(), http.StatusInternalServerError) + } + + return modTime, nil +} + func (a *App) MoveFile(oldPath, newPath string) *model.AppError { backend, err := a.FileBackend() if err != nil { diff --git a/app/import.go b/app/import.go index d8783d8e4b..dc5323856a 100644 --- a/app/import.go +++ b/app/import.go @@ -277,7 +277,7 @@ func (a *App) ListImports() ([]string, *model.AppError) { results := make([]string, 0, len(imports)) for i := 0; i < len(imports); i++ { filename := filepath.Base(imports[i]) - if !strings.HasSuffix(filename, incompleteUploadSuffix) { + if !strings.HasSuffix(filename, IncompleteUploadSuffix) { results = append(results, filename) } } diff --git a/app/opentracing/opentracing_layer.go b/app/opentracing/opentracing_layer.go index 0e264c9a6f..6741bfbca1 100644 --- a/app/opentracing/opentracing_layer.go +++ b/app/opentracing/opentracing_layer.go @@ -3689,6 +3689,28 @@ func (a *OpenTracingAppLayer) FileExists(path string) (bool, *model.AppError) { return resultVar0, resultVar1 } +func (a *OpenTracingAppLayer) FileModTime(path string) (time.Time, *model.AppError) { + origCtx := a.ctx + span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.FileModTime") + + a.ctx = newCtx + a.app.Srv().Store.SetContext(newCtx) + defer func() { + a.app.Srv().Store.SetContext(origCtx) + a.ctx = origCtx + }() + + defer span.Finish() + resultVar0, resultVar1 := a.app.FileModTime(path) + + if resultVar1 != nil { + span.LogFields(spanlog.Error(resultVar1)) + ext.Error.Set(span, true) + } + + return resultVar0, resultVar1 +} + func (a *OpenTracingAppLayer) FileReader(path string) (filesstore.ReadCloseSeeker, *model.AppError) { origCtx := a.ctx span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.FileReader") diff --git a/app/upload.go b/app/upload.go index e70976ca45..ea8e02e06f 100644 --- a/app/upload.go +++ b/app/upload.go @@ -19,7 +19,7 @@ import ( ) const minFirstPartSize = 5 * 1024 * 1024 // 5MB -const incompleteUploadSuffix = ".tmp" +const IncompleteUploadSuffix = ".tmp" func (a *App) runPluginsHook(info *model.FileInfo, file io.Reader) *model.AppError { pluginsEnvironment := a.GetPluginsEnvironment() @@ -111,7 +111,7 @@ func (a *App) CreateUploadSession(us *model.UploadSession) (*model.UploadSession if us.Type == model.UploadTypeAttachment { us.Path = now.Format("20060102") + "/teams/noteam/channels/" + us.ChannelId + "/users/" + us.UserId + "/" + us.Id + "/" + filepath.Base(us.Filename) } else if us.Type == model.UploadTypeImport { - us.Path = *a.Config().ImportSettings.Directory + "/" + us.Id + "_" + filepath.Base(us.Filename) + us.Path = filepath.Clean(*a.Config().ImportSettings.Directory) + "/" + us.Id + "_" + filepath.Base(us.Filename) } if err := us.IsValid(); err != nil { return nil, err @@ -194,7 +194,7 @@ func (a *App) UploadData(us *model.UploadSession, rd io.Reader) (*model.FileInfo uploadPath := us.Path if us.Type == model.UploadTypeImport { - uploadPath += incompleteUploadSuffix + uploadPath += IncompleteUploadSuffix } // make sure it's not possible to upload more data than what is expected. diff --git a/i18n/en.json b/i18n/en.json index 13f35f3d5b..28c91ebfca 100644 --- a/i18n/en.json +++ b/i18n/en.json @@ -1336,6 +1336,10 @@ "id": "api.file.file_exists.app_error", "translation": "Unable to check if the file exists." }, + { + "id": "api.file.file_mod_time.app_error", + "translation": "Unable to get last modification time for file." + }, { "id": "api.file.file_reader.app_error", "translation": "Unable to get a file reader." diff --git a/imports/placeholder.go b/imports/placeholder.go index 2eb8b1b15c..59c508f340 100644 --- a/imports/placeholder.go +++ b/imports/placeholder.go @@ -24,4 +24,7 @@ import ( // This is a placeholder so this package can be imported in Team Edition when it will be otherwise empty. _ "github.com/mattermost/mattermost-server/v5/jobs/import_process" + + // This is a placeholder so this package can be imported in Team Edition when it will be otherwise empty. + _ "github.com/mattermost/mattermost-server/v5/jobs/import_delete" ) diff --git a/jobs/import_delete/scheduler.go b/jobs/import_delete/scheduler.go new file mode 100644 index 0000000000..6982c5c352 --- /dev/null +++ b/jobs/import_delete/scheduler.go @@ -0,0 +1,51 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package import_delete + +import ( + "time" + + "github.com/mattermost/mattermost-server/v5/app" + "github.com/mattermost/mattermost-server/v5/model" +) + +const ( + jobName = "ImportDelete" + schedFrequency = 24 * time.Hour +) + +type Scheduler struct { + app *app.App +} + +func (i *ImportDeleteInterfaceImpl) MakeScheduler() model.Scheduler { + return &Scheduler{i.app} +} + +func (scheduler *Scheduler) Name() string { + return jobName + "Scheduler" +} + +func (scheduler *Scheduler) JobType() string { + return model.JOB_TYPE_IMPORT_DELETE +} + +func (scheduler *Scheduler) Enabled(cfg *model.Config) bool { + return *cfg.ImportSettings.Directory != "" && *cfg.ImportSettings.RetentionDays > 0 +} + +func (scheduler *Scheduler) NextScheduleTime(cfg *model.Config, now time.Time, pendingJobs bool, lastSuccessfulJob *model.Job) *time.Time { + nextTime := time.Now().Add(schedFrequency) + return &nextTime +} + +func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) { + data := map[string]string{} + + job, err := scheduler.app.Srv().Jobs.CreateJob(model.JOB_TYPE_IMPORT_DELETE, data) + if err != nil { + return nil, err + } + return job, nil +} diff --git a/jobs/import_delete/worker.go b/jobs/import_delete/worker.go new file mode 100644 index 0000000000..acee27ec39 --- /dev/null +++ b/jobs/import_delete/worker.go @@ -0,0 +1,172 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package import_delete + +import ( + "errors" + "path/filepath" + "time" + + "github.com/mattermost/mattermost-server/v5/app" + "github.com/mattermost/mattermost-server/v5/jobs" + tjobs "github.com/mattermost/mattermost-server/v5/jobs/interfaces" + "github.com/mattermost/mattermost-server/v5/mlog" + "github.com/mattermost/mattermost-server/v5/model" + "github.com/mattermost/mattermost-server/v5/store" +) + +func init() { + app.RegisterJobsImportDeleteInterface(func(a *app.App) tjobs.ImportDeleteInterface { + return &ImportDeleteInterfaceImpl{a} + }) +} + +type ImportDeleteInterfaceImpl struct { + app *app.App +} + +type ImportDeleteWorker struct { + name string + stopChan chan struct{} + stoppedChan chan struct{} + jobsChan chan model.Job + jobServer *jobs.JobServer + app *app.App +} + +func (i *ImportDeleteInterfaceImpl) MakeWorker() model.Worker { + return &ImportDeleteWorker{ + name: "ImportDelete", + stopChan: make(chan struct{}), + stoppedChan: make(chan struct{}), + jobsChan: make(chan model.Job), + jobServer: i.app.Srv().Jobs, + app: i.app, + } +} + +func (w *ImportDeleteWorker) JobChannel() chan<- model.Job { + return w.jobsChan +} + +func (w *ImportDeleteWorker) Run() { + mlog.Debug("Worker started", mlog.String("worker", w.name)) + + defer func() { + mlog.Debug("Worker finished", mlog.String("worker", w.name)) + close(w.stoppedChan) + }() + + for { + select { + case <-w.stopChan: + mlog.Debug("Worker received stop signal", mlog.String("worker", w.name)) + return + case job := <-w.jobsChan: + mlog.Debug("Worker received a new candidate job.", mlog.String("worker", w.name)) + w.doJob(&job) + } + } +} + +func (w *ImportDeleteWorker) Stop() { + mlog.Debug("Worker stopping", mlog.String("worker", w.name)) + close(w.stopChan) + <-w.stoppedChan +} + +func (w *ImportDeleteWorker) doJob(job *model.Job) { + if claimed, err := w.jobServer.ClaimJob(job); err != nil { + mlog.Warn("Worker experienced an error while trying to claim job", + mlog.String("worker", w.name), + mlog.String("job_id", job.Id), + mlog.String("error", err.Error())) + return + } else if !claimed { + return + } + + importPath := *w.app.Config().ImportSettings.Directory + retentionTime := time.Duration(*w.app.Config().ImportSettings.RetentionDays) * 24 * time.Hour + imports, appErr := w.app.ListDirectory(importPath) + if appErr != nil { + w.setJobError(job, appErr) + return + } + + var hasErrs bool + for i := range imports { + filename := filepath.Base(imports[i]) + modTime, appErr := w.app.FileModTime(filepath.Join(importPath, filename)) + if appErr != nil { + mlog.Debug("Worker: Failed to get file modification time", + mlog.Err(appErr), mlog.String("import", imports[i])) + hasErrs = true + continue + } + + if time.Now().After(modTime.Add(retentionTime)) { + // expected format if uploaded through the API is + // ${uploadID}_${filename}${app.IncompleteUploadSuffix} + minLen := 26 + 1 + len(app.IncompleteUploadSuffix) + + // check if it's an incomplete upload and attempt to delete its session. + if len(filename) > minLen && filepath.Ext(filename) == app.IncompleteUploadSuffix { + uploadID := filename[:26] + if storeErr := w.app.Srv().Store.UploadSession().Delete(uploadID); storeErr != nil { + mlog.Debug("Worker: Failed to delete UploadSession", + mlog.Err(storeErr), mlog.String("upload_id", uploadID)) + hasErrs = true + continue + } + } else { + // check if fileinfo exists and if so delete it. + filePath := filepath.Join(imports[i]) + info, storeErr := w.app.Srv().Store.FileInfo().GetByPath(filePath) + var nfErr *store.ErrNotFound + if storeErr != nil && !errors.As(storeErr, &nfErr) { + mlog.Debug("Worker: Failed to get FileInfo", + mlog.Err(storeErr), mlog.String("path", filePath)) + hasErrs = true + continue + } else if storeErr == nil { + if storeErr = w.app.Srv().Store.FileInfo().PermanentDelete(info.Id); storeErr != nil { + mlog.Debug("Worker: Failed to delete FileInfo", + mlog.Err(storeErr), mlog.String("file_id", info.Id)) + hasErrs = true + continue + } + } + } + + // remove file data from storage. + if appErr := w.app.RemoveFile(imports[i]); appErr != nil { + mlog.Debug("Worker: Failed to remove file", + mlog.Err(appErr), mlog.String("import", imports[i])) + hasErrs = true + continue + } + } + } + + if hasErrs { + mlog.Warn("Worker: errors occurred") + } + + mlog.Info("Worker: Job is complete", mlog.String("worker", w.name), mlog.String("job_id", job.Id)) + w.setJobSuccess(job) +} + +func (w *ImportDeleteWorker) setJobSuccess(job *model.Job) { + if err := w.app.Srv().Jobs.SetJobSuccess(job); err != nil { + mlog.Error("Worker: Failed to set success for job", mlog.String("worker", w.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error())) + w.setJobError(job, err) + } +} + +func (w *ImportDeleteWorker) setJobError(job *model.Job, appError *model.AppError) { + if err := w.app.Srv().Jobs.SetJobError(job, appError); err != nil { + mlog.Error("Worker: Failed to set job error", mlog.String("worker", w.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error())) + } +} diff --git a/jobs/interfaces/import_delete_interface.go b/jobs/interfaces/import_delete_interface.go new file mode 100644 index 0000000000..8506ec8ff2 --- /dev/null +++ b/jobs/interfaces/import_delete_interface.go @@ -0,0 +1,11 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package interfaces + +import "github.com/mattermost/mattermost-server/v5/model" + +type ImportDeleteInterface interface { + MakeWorker() model.Worker + MakeScheduler() model.Scheduler +} diff --git a/jobs/jobs_watcher.go b/jobs/jobs_watcher.go index b8c507cd39..2ca31e0096 100644 --- a/jobs/jobs_watcher.go +++ b/jobs/jobs_watcher.go @@ -156,6 +156,13 @@ func (watcher *Watcher) PollAndNotify() { default: } } + } else if job.Type == model.JOB_TYPE_IMPORT_DELETE { + if watcher.workers.ImportDelete != nil { + select { + case watcher.workers.ImportDelete.JobChannel() <- *job: + default: + } + } } else if job.Type == model.JOB_TYPE_CLOUD { if watcher.workers.Cloud != nil { select { diff --git a/jobs/schedulers.go b/jobs/schedulers.go index c1820a78bc..9a2f479634 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -69,6 +69,7 @@ func (srv *JobServer) InitSchedulers() *Schedulers { if activeUsersInterface := srv.ActiveUsers; activeUsersInterface != nil { schedulers.schedulers = append(schedulers.schedulers, activeUsersInterface.MakeScheduler()) } + if productNoticesInterface := srv.ProductNotices; productNoticesInterface != nil { schedulers.schedulers = append(schedulers.schedulers, productNoticesInterface.MakeScheduler()) } @@ -77,6 +78,10 @@ func (srv *JobServer) InitSchedulers() *Schedulers { schedulers.schedulers = append(schedulers.schedulers, cloudInterface.MakeScheduler()) } + if importDeleteInterface := srv.ImportDelete; importDeleteInterface != nil { + schedulers.schedulers = append(schedulers.schedulers, importDeleteInterface.MakeScheduler()) + } + schedulers.nextRunTimes = make([]*time.Time, len(schedulers.schedulers)) return schedulers } diff --git a/jobs/server.go b/jobs/server.go index b89ea595f2..4b30a6a88e 100644 --- a/jobs/server.go +++ b/jobs/server.go @@ -31,6 +31,7 @@ type JobServer struct { ProductNotices tjobs.ProductNoticesJobInterface ActiveUsers tjobs.ActiveUsersJobInterface ImportProcess tjobs.ImportProcessInterface + ImportDelete tjobs.ImportDeleteInterface Cloud ejobs.CloudJobInterface } diff --git a/jobs/workers.go b/jobs/workers.go index 0efa5802eb..acce58f8ac 100644 --- a/jobs/workers.go +++ b/jobs/workers.go @@ -28,6 +28,7 @@ type Workers struct { ProductNotices model.Worker ActiveUsers model.Worker ImportProcess model.Worker + ImportDelete model.Worker Cloud model.Worker listenerId string @@ -87,6 +88,10 @@ func (srv *JobServer) InitWorkers() *Workers { workers.ImportProcess = importProcessInterface.MakeWorker() } + if importDeleteInterface := srv.ImportDelete; importDeleteInterface != nil { + workers.ImportDelete = importDeleteInterface.MakeWorker() + } + if cloudInterface := srv.Cloud; cloudInterface != nil { workers.Cloud = cloudInterface.MakeWorker() } @@ -146,6 +151,10 @@ func (workers *Workers) Start() *Workers { go workers.ImportProcess.Run() } + if workers.ImportDelete != nil { + go workers.ImportDelete.Run() + } + if workers.Cloud != nil { go workers.Cloud.Run() } @@ -263,6 +272,10 @@ func (workers *Workers) Stop() *Workers { workers.ImportProcess.Stop() } + if workers.ImportDelete != nil { + workers.ImportDelete.Stop() + } + if workers.Cloud != nil { workers.Cloud.Stop() } diff --git a/model/job.go b/model/job.go index 0c006ea44c..b2cdaba122 100644 --- a/model/job.go +++ b/model/job.go @@ -23,6 +23,7 @@ const ( JOB_TYPE_PRODUCT_NOTICES = "product_notices" JOB_TYPE_ACTIVE_USERS = "active_users" JOB_TYPE_IMPORT_PROCESS = "import_process" + JOB_TYPE_IMPORT_DELETE = "import_delete" JOB_TYPE_CLOUD = "cloud" JOB_STATUS_PENDING = "pending" @@ -68,6 +69,7 @@ func (j *Job) IsValid() *AppError { case JOB_TYPE_EXPIRY_NOTIFY: case JOB_TYPE_ACTIVE_USERS: case JOB_TYPE_IMPORT_PROCESS: + case JOB_TYPE_IMPORT_DELETE: case JOB_TYPE_CLOUD: default: return NewAppError("Job.IsValid", "model.job.is_valid.type.app_error", nil, "id="+j.Id, http.StatusBadRequest) diff --git a/services/filesstore/filesstore.go b/services/filesstore/filesstore.go index c2726bba4b..251b36957f 100644 --- a/services/filesstore/filesstore.go +++ b/services/filesstore/filesstore.go @@ -5,6 +5,7 @@ package filesstore import ( "io" + "time" "github.com/pkg/errors" @@ -28,6 +29,7 @@ type FileBackend interface { WriteFile(fr io.Reader, path string) (int64, error) AppendFile(fr io.Reader, path string) (int64, error) RemoveFile(path string) error + FileModTime(path string) (time.Time, error) ListDirectory(path string) ([]string, error) RemoveDirectory(path string) error diff --git a/services/filesstore/filesstore_test.go b/services/filesstore/filesstore_test.go index 0dc4077f1d..3f500551f4 100644 --- a/services/filesstore/filesstore_test.go +++ b/services/filesstore/filesstore_test.go @@ -10,6 +10,7 @@ import ( "math/rand" "os" "testing" + "time" "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" @@ -388,6 +389,42 @@ func (s *FileBackendTestSuite) TestFileSize() { }) } +func (s *FileBackendTestSuite) TestFileModTime() { + s.Run("nonexistent file", func() { + modTime, err := s.backend.FileModTime("tests/nonexistentfile") + s.NotNil(err) + s.Empty(modTime) + }) + + s.Run("valid file", func() { + path := "tests/" + model.NewId() + data := []byte("some data") + + written, err := s.backend.WriteFile(bytes.NewReader(data), path) + s.Nil(err) + s.EqualValues(len(data), written) + defer s.backend.RemoveFile(path) + + modTime, err := s.backend.FileModTime(path) + s.Nil(err) + s.NotEmpty(modTime) + + // We wait 1 second so that the times will differ enough to be testable. + time.Sleep(1 * time.Second) + + path2 := "tests/" + model.NewId() + written, err = s.backend.WriteFile(bytes.NewReader(data), path2) + s.Nil(err) + s.EqualValues(len(data), written) + defer s.backend.RemoveFile(path2) + + modTime2, err := s.backend.FileModTime(path2) + s.Nil(err) + s.NotEmpty(modTime2) + s.True(modTime2.After(modTime)) + }) +} + func BenchmarkS3WriteFile(b *testing.B) { utils.TranslationsPreInit() diff --git a/services/filesstore/localstore.go b/services/filesstore/localstore.go index 14ed45c956..76c68d8ad0 100644 --- a/services/filesstore/localstore.go +++ b/services/filesstore/localstore.go @@ -9,6 +9,7 @@ import ( "io/ioutil" "os" "path/filepath" + "time" "github.com/pkg/errors" @@ -71,6 +72,14 @@ func (b *LocalFileBackend) FileSize(path string) (int64, error) { return info.Size(), nil } +func (b *LocalFileBackend) FileModTime(path string) (time.Time, error) { + info, err := os.Stat(filepath.Join(b.directory, path)) + if err != nil { + return time.Time{}, errors.Wrapf(err, "unable to get modification time for file %s", path) + } + return info.ModTime(), nil +} + func (b *LocalFileBackend) CopyFile(oldPath, newPath string) error { if err := utils.CopyFile(filepath.Join(b.directory, oldPath), filepath.Join(b.directory, newPath)); err != nil { return errors.Wrapf(err, "unable to copy file from %s to %s", oldPath, newPath) diff --git a/services/filesstore/mocks/FileBackend.go b/services/filesstore/mocks/FileBackend.go index 454bb1dab0..33dcf49e79 100644 --- a/services/filesstore/mocks/FileBackend.go +++ b/services/filesstore/mocks/FileBackend.go @@ -10,6 +10,8 @@ import ( filesstore "github.com/mattermost/mattermost-server/v5/services/filesstore" mock "github.com/stretchr/testify/mock" + + time "time" ) // FileBackend is an autogenerated mock type for the FileBackend type @@ -73,6 +75,27 @@ func (_m *FileBackend) FileExists(path string) (bool, error) { return r0, r1 } +// FileModTime provides a mock function with given fields: path +func (_m *FileBackend) FileModTime(path string) (time.Time, error) { + ret := _m.Called(path) + + var r0 time.Time + if rf, ok := ret.Get(0).(func(string) time.Time); ok { + r0 = rf(path) + } else { + r0 = ret.Get(0).(time.Time) + } + + var r1 error + if rf, ok := ret.Get(1).(func(string) error); ok { + r1 = rf(path) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + // FileSize provides a mock function with given fields: path func (_m *FileBackend) FileSize(path string) (int64, error) { ret := _m.Called(path) diff --git a/services/filesstore/s3store.go b/services/filesstore/s3store.go index 9430879183..2b00bc6e11 100644 --- a/services/filesstore/s3store.go +++ b/services/filesstore/s3store.go @@ -10,6 +10,7 @@ import ( "os" "path/filepath" "strings" + "time" s3 "github.com/minio/minio-go/v7" "github.com/minio/minio-go/v7/pkg/credentials" @@ -205,6 +206,17 @@ func (b *S3FileBackend) FileSize(path string) (int64, error) { return info.Size, nil } +func (b *S3FileBackend) FileModTime(path string) (time.Time, error) { + path = filepath.Join(b.pathPrefix, path) + + info, err := b.client.StatObject(context.Background(), b.bucket, path, s3.StatObjectOptions{}) + if err != nil { + return time.Time{}, errors.Wrapf(err, "unable to get modification time for file %s", path) + } + + return info.LastModified, nil +} + func (b *S3FileBackend) CopyFile(oldPath, newPath string) error { oldPath = filepath.Join(b.pathPrefix, oldPath) newPath = filepath.Join(b.pathPrefix, newPath)