diff --git a/app/server.go b/app/server.go index 56688be735..7bf6de85ba 100644 --- a/app/server.go +++ b/app/server.go @@ -760,6 +760,9 @@ func (s *Server) runJobs() { s.Go(func() { runSessionCleanupJob(s) }) + s.Go(func() { + runJobsCleanupJob(s) + }) s.Go(func() { runTokenCleanupJob(s) }) @@ -1504,6 +1507,13 @@ func runSessionCleanupJob(s *Server) { }, time.Hour*24) } +func runJobsCleanupJob(s *Server) { + doJobsCleanup(s) + model.CreateRecurringTask("Job Cleanup", func() { + doJobsCleanup(s) + }, time.Hour*24) +} + func (s *Server) runLicenseExpirationCheckJob() { s.doLicenseExpirationCheck() model.CreateRecurringTask("License Expiration Check", func() { @@ -1549,6 +1559,7 @@ func doCommandWebhookCleanup(s *Server) { const ( sessionsCleanupBatchSize = 1000 + jobsCleanupBatchSize = 1000 ) func doSessionCleanup(s *Server) { @@ -1559,6 +1570,20 @@ func doSessionCleanup(s *Server) { } } +func doJobsCleanup(s *Server) { + if *s.Config().JobSettings.CleanupJobsThresholdDays < 0 { + return + } + mlog.Debug("Cleaning up jobs store.") + + dur := time.Duration(*s.Config().JobSettings.CleanupJobsThresholdDays) * time.Hour * 24 + expiry := model.GetMillisForTime(time.Now().Add(-dur)) + err := s.Store.Job().Cleanup(expiry, jobsCleanupBatchSize) + if err != nil { + mlog.Warn("Error while cleaning up jobs", mlog.Err(err)) + } +} + func doCheckAdminSupportStatus(a *App, c *request.Context) { isE0Edition := model.BuildEnterpriseReady == "true" diff --git a/model/config.go b/model/config.go index 059fc21c8c..5699e068d3 100644 --- a/model/config.go +++ b/model/config.go @@ -2615,8 +2615,9 @@ func (s *DataRetentionSettings) SetDefaults() { } type JobSettings struct { - RunJobs *bool `access:"write_restrictable,cloud_restrictable"` - RunScheduler *bool `access:"write_restrictable,cloud_restrictable"` + RunJobs *bool `access:"write_restrictable,cloud_restrictable"` + RunScheduler *bool `access:"write_restrictable,cloud_restrictable"` + CleanupJobsThresholdDays *int `access:"write_restrictable,cloud_restrictable"` } func (s *JobSettings) SetDefaults() { @@ -2627,6 +2628,10 @@ func (s *JobSettings) SetDefaults() { if s.RunScheduler == nil { s.RunScheduler = NewBool(true) } + + if s.CleanupJobsThresholdDays == nil { + s.CleanupJobsThresholdDays = NewInt(-1) + } } type CloudSettings struct { diff --git a/store/opentracinglayer/opentracinglayer.go b/store/opentracinglayer/opentracinglayer.go index 48d847c045..cf038cf5a4 100644 --- a/store/opentracinglayer/opentracinglayer.go +++ b/store/opentracinglayer/opentracinglayer.go @@ -4177,6 +4177,24 @@ func (s *OpenTracingLayerGroupStore) UpsertMember(groupID string, userID string) return result, err } +func (s *OpenTracingLayerJobStore) Cleanup(expiryTime int64, batchSize int) error { + origCtx := s.Root.Store.Context() + span, newCtx := tracing.StartSpanWithParentByContext(s.Root.Store.Context(), "JobStore.Cleanup") + s.Root.Store.SetContext(newCtx) + defer func() { + s.Root.Store.SetContext(origCtx) + }() + + defer span.Finish() + err := s.JobStore.Cleanup(expiryTime, batchSize) + if err != nil { + span.LogFields(spanlog.Error(err)) + ext.Error.Set(span, true) + } + + return err +} + func (s *OpenTracingLayerJobStore) Delete(id string) (string, error) { origCtx := s.Root.Store.Context() span, newCtx := tracing.StartSpanWithParentByContext(s.Root.Store.Context(), "JobStore.Delete") diff --git a/store/retrylayer/retrylayer.go b/store/retrylayer/retrylayer.go index 69f4b2340f..418c38b197 100644 --- a/store/retrylayer/retrylayer.go +++ b/store/retrylayer/retrylayer.go @@ -4506,6 +4506,26 @@ func (s *RetryLayerGroupStore) UpsertMember(groupID string, userID string) (*mod } +func (s *RetryLayerJobStore) Cleanup(expiryTime int64, batchSize int) error { + + tries := 0 + for { + err := s.JobStore.Cleanup(expiryTime, batchSize) + if err == nil { + return nil + } + if !isRepeatableError(err) { + return err + } + tries++ + if tries >= 3 { + err = errors.Wrap(err, "giving up after 3 consecutive repeatable transaction failures") + return err + } + } + +} + func (s *RetryLayerJobStore) Delete(id string) (string, error) { tries := 0 diff --git a/store/sqlstore/job_store.go b/store/sqlstore/job_store.go index 12ffad2654..85bfd91fb1 100644 --- a/store/sqlstore/job_store.go +++ b/store/sqlstore/job_store.go @@ -8,6 +8,7 @@ import ( "encoding/json" "fmt" "strings" + "time" sq "github.com/Masterminds/squirrel" "github.com/mattermost/gorp" @@ -17,6 +18,10 @@ import ( "github.com/mattermost/mattermost-server/v6/store" ) +const ( + jobsCleanupDelay = 100 * time.Millisecond +) + type SqlJobStore struct { *SqlStore } @@ -286,3 +291,31 @@ func (jss SqlJobStore) Delete(id string) (string, error) { } return id, nil } + +func (jss SqlJobStore) Cleanup(expiryTime int64, batchSize int) error { + var query string + if jss.DriverName() == model.DatabaseDriverPostgres { + query = "DELETE FROM Jobs WHERE Id IN (SELECT Id FROM Jobs WHERE CreateAt < ? AND (Status != ? AND Status != ?) ORDER BY CreateAt ASC LIMIT ?)" + } else { + query = "DELETE FROM Jobs WHERE CreateAt < ? AND (Status != ? AND Status != ?) ORDER BY CreateAt ASC LIMIT ?" + } + + var rowsAffected int64 = 1 + + for rowsAffected > 0 { + sqlResult, err := jss.GetMasterX().Exec(query, + expiryTime, model.JobStatusInProgress, model.JobStatusPending, batchSize) + if err != nil { + return errors.Wrap(err, "unable to delete jobs") + } + var rowErr error + rowsAffected, rowErr = sqlResult.RowsAffected() + if rowErr != nil { + return errors.Wrap(err, "unable to delete jobs") + } + + time.Sleep(jobsCleanupDelay) + } + + return nil +} diff --git a/store/store.go b/store/store.go index be4cdadfc2..67828b6942 100644 --- a/store/store.go +++ b/store/store.go @@ -681,6 +681,7 @@ type JobStore interface { GetNewestJobByStatusesAndType(statuses []string, jobType string) (*model.Job, error) GetCountByStatusAndType(status string, jobType string) (int64, error) Delete(id string) (string, error) + Cleanup(expiryTime int64, batchSize int) error } type UserAccessTokenStore interface { diff --git a/store/storetest/job_store.go b/store/storetest/job_store.go index c7ce43cd27..5ac9d4051d 100644 --- a/store/storetest/job_store.go +++ b/store/storetest/job_store.go @@ -29,6 +29,7 @@ func TestJobStore(t *testing.T, ss store.Store) { t.Run("JobUpdateOptimistically", func(t *testing.T) { testJobUpdateOptimistically(t, ss) }) t.Run("JobUpdateStatusUpdateStatusOptimistically", func(t *testing.T) { testJobUpdateStatusUpdateStatusOptimistically(t, ss) }) t.Run("JobDelete", func(t *testing.T) { testJobDelete(t, ss) }) + t.Run("JobCleanup", func(t *testing.T) { testJobCleanup(t, ss) }) } func testJobSaveGet(t *testing.T, ss store.Store) { @@ -552,3 +553,43 @@ func testJobDelete(t *testing.T, ss store.Store) { _, err = ss.Job().Delete(job.Id) assert.NoError(t, err) } + +func testJobCleanup(t *testing.T, ss store.Store) { + now := model.GetMillis() + ids := make([]string, 0, 10) + for i := 0; i < 10; i++ { + job, err := ss.Job().Save(&model.Job{ + Id: model.NewId(), + CreateAt: now - int64(i), + Status: model.JobStatusPending, + }) + require.NoError(t, err) + ids = append(ids, job.Id) + defer ss.Job().Delete(job.Id) + } + + jobs, err := ss.Job().GetAllByStatus(model.JobStatusPending) + require.NoError(t, err) + assert.Len(t, jobs, 10) + + err = ss.Job().Cleanup(now+1, 5) + require.NoError(t, err) + + // Should not clean up pending jobs + jobs, err = ss.Job().GetAllByStatus(model.JobStatusPending) + require.NoError(t, err) + assert.Len(t, jobs, 10) + + for _, id := range ids { + _, err = ss.Job().UpdateStatus(id, model.JobStatusSuccess) + require.NoError(t, err) + } + + err = ss.Job().Cleanup(now+1, 5) + require.NoError(t, err) + + // Should clean up now + jobs, err = ss.Job().GetAllByStatus(model.JobStatusSuccess) + require.NoError(t, err) + assert.Len(t, jobs, 0) +} diff --git a/store/storetest/mocks/JobStore.go b/store/storetest/mocks/JobStore.go index 83c0f27c46..341841da49 100644 --- a/store/storetest/mocks/JobStore.go +++ b/store/storetest/mocks/JobStore.go @@ -14,6 +14,20 @@ type JobStore struct { mock.Mock } +// Cleanup provides a mock function with given fields: expiryTime, batchSize +func (_m *JobStore) Cleanup(expiryTime int64, batchSize int) error { + ret := _m.Called(expiryTime, batchSize) + + var r0 error + if rf, ok := ret.Get(0).(func(int64, int) error); ok { + r0 = rf(expiryTime, batchSize) + } else { + r0 = ret.Error(0) + } + + return r0 +} + // Delete provides a mock function with given fields: id func (_m *JobStore) Delete(id string) (string, error) { ret := _m.Called(id) diff --git a/store/timerlayer/timerlayer.go b/store/timerlayer/timerlayer.go index 40f80990e5..12e231387f 100644 --- a/store/timerlayer/timerlayer.go +++ b/store/timerlayer/timerlayer.go @@ -3803,6 +3803,22 @@ func (s *TimerLayerGroupStore) UpsertMember(groupID string, userID string) (*mod return result, err } +func (s *TimerLayerJobStore) Cleanup(expiryTime int64, batchSize int) error { + start := timemodule.Now() + + err := s.JobStore.Cleanup(expiryTime, batchSize) + + elapsed := float64(timemodule.Since(start)) / float64(timemodule.Second) + if s.Root.Metrics != nil { + success := "false" + if err == nil { + success = "true" + } + s.Root.Metrics.ObserveStoreMethodDuration("JobStore.Cleanup", success, elapsed) + } + return err +} + func (s *TimerLayerJobStore) Delete(id string) (string, error) { start := timemodule.Now()