MM-17508: Modifying config files causes compliance exports to run twice (#13053)
* Check if leader before setting schedule. * add check to normal run as well
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
988b8b74c5
Коммит
2bcb4d9913
@@ -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 {
|
||||
|
||||
105
jobs/schedulers_test.go
Обычный файл
105
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()
|
||||
|
||||
}
|
||||
Ссылка в новой задаче
Block a user