diff --git a/jobs/schedulers.go b/jobs/schedulers.go index 52df4eebea..704fbb5759 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -20,6 +20,7 @@ type Schedulers struct { listenerId string startOnce sync.Once jobs *JobServer + isLeader bool schedulers []model.Scheduler nextRunTimes []*time.Time @@ -34,6 +35,7 @@ func (srv *JobServer) InitSchedulers() *Schedulers { configChanged: make(chan *model.Config), clusterLeaderChanged: make(chan bool), jobs: srv, + isLeader: true, } if srv.DataRetentionJob != nil { @@ -103,7 +105,7 @@ func (schedulers *Schedulers) Start() *Schedulers { if time.Now().After(*nextTime) { scheduler := schedulers.schedulers[idx] if scheduler != nil { - if scheduler.Enabled(cfg) { + if schedulers.isLeader && scheduler.Enabled(cfg) { if _, err := schedulers.scheduleJob(cfg, scheduler); err != nil { mlog.Error("Failed to schedule job", mlog.String("scheduler", scheduler.Name()), mlog.Err(err)) } else { @@ -115,7 +117,7 @@ func (schedulers *Schedulers) Start() *Schedulers { } case newCfg := <-schedulers.configChanged: for idx, scheduler := range schedulers.schedulers { - if !scheduler.Enabled(newCfg) { + if !schedulers.isLeader || !scheduler.Enabled(newCfg) { schedulers.nextRunTimes[idx] = nil } else { schedulers.setNextRunTime(newCfg, idx, now, false) @@ -123,6 +125,7 @@ func (schedulers *Schedulers) Start() *Schedulers { } case isLeader := <-schedulers.clusterLeaderChanged: for idx := range schedulers.schedulers { + schedulers.isLeader = isLeader if !isLeader { schedulers.nextRunTimes[idx] = nil } else { diff --git a/jobs/schedulers_test.go b/jobs/schedulers_test.go new file mode 100644 index 0000000000..35dd4fc5f8 --- /dev/null +++ b/jobs/schedulers_test.go @@ -0,0 +1,105 @@ +// Copyright (c) 2017-present Mattermost, Inc. All Rights Reserved. +// See License.txt for license information. +package jobs + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" + + "github.com/mattermost/mattermost-server/einterfaces/mocks" + "github.com/mattermost/mattermost-server/plugin/plugintest/mock" + + "github.com/mattermost/mattermost-server/model" + "github.com/mattermost/mattermost-server/store/storetest" + "github.com/mattermost/mattermost-server/utils/testutils" +) + +type MockScheduler struct { + mock.Mock +} + +func (scheduler *MockScheduler) Enabled(cfg *model.Config) bool { + return true +} + +func (scheduler *MockScheduler) Name() string { + return "MockScheduler" +} + +func (scheduler *MockScheduler) JobType() string { + return model.JOB_TYPE_DATA_RETENTION +} + +func (scheduler *MockScheduler) NextScheduleTime(cfg *model.Config, now time.Time, pendingJobs bool, lastSuccessfulJob *model.Job) *time.Time { + nextTime := time.Now().Add(60 * time.Second) + return &nextTime +} + +func (scheduler *MockScheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) { + return nil, nil +} + +func TestScheduler(t *testing.T) { + mockStore := &storetest.Store{} + defer mockStore.AssertExpectations(t) + + job := &model.Job{ + Id: model.NewId(), + CreateAt: model.GetMillis(), + Status: model.JOB_STATUS_PENDING, + Type: model.JOB_TYPE_MESSAGE_EXPORT, + } + // mock job store doesn't return a previously successful job, forcing fallback to config + mockStore.JobStore.On("GetNewestJobByStatusAndType", mock.AnythingOfType("string"), mock.AnythingOfType("string")).Return(job, nil) + mockStore.JobStore.On("GetCountByStatusAndType", mock.AnythingOfType("string"), mock.AnythingOfType("string")).Return(int64(1), nil) + + jobServer := &JobServer{ + Store: mockStore, + ConfigService: &testutils.StaticConfigService{ + Cfg: &model.Config{ + // mock config + DataRetentionSettings: *&model.DataRetentionSettings{ + EnableMessageDeletion: model.NewBool(true), + }, + MessageExportSettings: *&model.MessageExportSettings{ + EnableExport: model.NewBool(true), + }, + }, + }, + } + + jobInterface := new(mocks.DataRetentionJobInterface) + jobInterface.On("MakeScheduler").Return(new(MockScheduler)) + jobServer.DataRetentionJob = jobInterface + + exportInterface := new(mocks.MessageExportJobInterface) + exportInterface.On("MakeScheduler").Return(new(MockScheduler)) + jobServer.MessageExportJob = exportInterface + + schedulers := jobServer.InitSchedulers() + schedulers.Start() + time.Sleep(1 * time.Second) + + // They should be all on here + for _, element := range schedulers.nextRunTimes { + assert.NotNil(t, element) + } + + schedulers.HandleClusterLeaderChange(false) + time.Sleep(1 * time.Second) + // They should be turned off + for _, element := range schedulers.nextRunTimes { + assert.Nil(t, element) + } + + // After running a config change, they should stay off + schedulers.handleConfigChange(nil, nil) + for _, element := range schedulers.nextRunTimes { + assert.Nil(t, element) + } + + schedulers.Stop() + +}