diff --git a/app/helper_test.go b/app/helper_test.go index 90d6893cd4..58063b9a10 100644 --- a/app/helper_test.go +++ b/app/helper_test.go @@ -134,8 +134,6 @@ func setupTestHelper(dbStore store.Store, enterprise bool, includeCacheLayer boo if enterprise { th.App.Srv().SetLicense(model.NewTestLicense()) - th.App.Srv().Jobs.InitWorkers() - th.App.Srv().Jobs.InitSchedulers() } else { th.App.Srv().SetLicense(nil) } diff --git a/app/license.go b/app/license.go index f6e9bafc10..779728f914 100644 --- a/app/license.go +++ b/app/license.go @@ -177,6 +177,34 @@ func (s *Server) SaveLicense(licenseBytes []byte) (*model.License, *model.AppErr return nil, model.NewAppError("addLicense", model.ExpiredLicenseError, nil, "", http.StatusBadRequest) } + if *s.Config().JobSettings.RunJobs && s.Jobs != nil { + if err := s.Jobs.StopWorkers(); err != nil && !errors.Is(err, jobs.ErrWorkersNotRunning) { + mlog.Warn("Stopping job server workers failed", mlog.Err(err)) + } + } + + if *s.Config().JobSettings.RunScheduler && s.Jobs != nil { + if err := s.Jobs.StopSchedulers(); err != nil && !errors.Is(err, jobs.ErrSchedulersNotRunning) { + mlog.Error("Stopping job server schedulers failed", mlog.Err(err)) + } + } + + defer func() { + // restart job server workers - this handles the edge case where a license file is uploaded, but the job server + // doesn't start until the server is restarted, which prevents the 'run job now' buttons in system console from + // functioning as expected + if *s.Config().JobSettings.RunJobs && s.Jobs != nil { + if err := s.Jobs.StartWorkers(); err != nil { + mlog.Error("Starting job server workers failed", mlog.Err(err)) + } + } + if *s.Config().JobSettings.RunScheduler && s.Jobs != nil { + if err := s.Jobs.StartSchedulers(); err != nil && !errors.Is(err, jobs.ErrSchedulersRunning) { + mlog.Error("Starting job server schedulers failed", mlog.Err(err)) + } + } + }() + if ok := s.SetLicense(&license); !ok { return nil, model.NewAppError("addLicense", model.ExpiredLicenseError, nil, "", http.StatusBadRequest) } @@ -208,25 +236,6 @@ func (s *Server) SaveLicense(licenseBytes []byte) (*model.License, *model.AppErr s.ReloadConfig() s.InvalidateAllCaches() - // restart job server workers - this handles the edge case where a license file is uploaded, but the job server - // doesn't start until the server is restarted, which prevents the 'run job now' buttons in system console from - // functioning as expected - if *s.Config().JobSettings.RunJobs && s.Jobs != nil { - if err := s.Jobs.StopWorkers(); err != nil && !errors.Is(err, jobs.ErrWorkersNotRunning) { - mlog.Warn("Stopping job server workers failed", mlog.Err(err)) - } - if err := s.Jobs.InitWorkers(); err != nil { - mlog.Error("Initializing job server workers failed", mlog.Err(err)) - } else if err := s.Jobs.StartWorkers(); err != nil { - mlog.Error("Starting job server workers failed", mlog.Err(err)) - } - } - if *s.Config().JobSettings.RunScheduler && s.Jobs != nil { - if err := s.Jobs.StartSchedulers(); err != nil && !errors.Is(err, jobs.ErrSchedulersRunning) { - mlog.Error("Starting job server schedulers failed", mlog.Err(err)) - } - } - return &license, nil } diff --git a/app/server.go b/app/server.go index 90f58f6fa9..018168b4f6 100644 --- a/app/server.go +++ b/app/server.go @@ -1879,8 +1879,6 @@ func (ch *Channels) ClientConfigHash() string { func (s *Server) initJobs() { s.Jobs = jobs.NewJobServer(s, s.Store, s.Metrics) - s.Jobs.InitWorkers() - s.Jobs.InitSchedulers() if jobsDataRetentionJobInterface != nil { builder := jobsDataRetentionJobInterface(s) diff --git a/app/slashcommands/helper_test.go b/app/slashcommands/helper_test.go index 66d1777d2c..56ba662f7f 100644 --- a/app/slashcommands/helper_test.go +++ b/app/slashcommands/helper_test.go @@ -125,9 +125,13 @@ func setupTestHelper(dbStore store.Store, enterprise bool, includeCacheLayer boo }) if enterprise { + th.App.Srv().Jobs.StopWorkers() + th.App.Srv().Jobs.StopSchedulers() + th.App.Srv().SetLicense(model.NewTestLicense()) - th.App.Srv().Jobs.InitWorkers() - th.App.Srv().Jobs.InitSchedulers() + + th.App.Srv().Jobs.StartWorkers() + th.App.Srv().Jobs.StartSchedulers() } else { th.App.Srv().SetLicense(nil) } diff --git a/jobs/jobs_test.go b/jobs/jobs_test.go index fb17d45b6c..bb28db69a8 100644 --- a/jobs/jobs_test.go +++ b/jobs/jobs_test.go @@ -28,7 +28,11 @@ func makeJobServer(t *testing.T) (*JobServer, *storetest.Store, *mocks.MetricsIn mockMetrics.AssertExpectations(t) }) - jobServer := NewJobServer(configService, mockStore, mockMetrics) + jobServer := &JobServer{ + ConfigService: configService, + Store: mockStore, + metrics: mockMetrics, + } return jobServer, mockStore, mockMetrics } diff --git a/jobs/jobs_watcher.go b/jobs/jobs_watcher.go index d8fcf71d55..e0f57a5cf6 100644 --- a/jobs/jobs_watcher.go +++ b/jobs/jobs_watcher.go @@ -26,8 +26,6 @@ type Watcher struct { func (srv *JobServer) MakeWatcher(workers *Workers, pollingInterval int) *Watcher { return &Watcher{ - stop: make(chan struct{}), - stopped: make(chan struct{}), pollingInterval: pollingInterval, workers: workers, srv: srv, @@ -36,7 +34,8 @@ func (srv *JobServer) MakeWatcher(workers *Workers, pollingInterval int) *Watche func (watcher *Watcher) Start() { mlog.Debug("Watcher Started") - + watcher.stop = make(chan struct{}) + watcher.stopped = make(chan struct{}) // Delay for some random number of milliseconds before starting to ensure that multiple // instances of the jobserver don't poll at a time too close to each other. rand.Seed(time.Now().UTC().UnixNano()) diff --git a/jobs/schedulers.go b/jobs/schedulers.go index 01aca21baa..f692b08451 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -32,30 +32,6 @@ var ( ErrSchedulersUninitialized = errors.New("job schedulers are not initialized") ) -func (srv *JobServer) InitSchedulers() error { - srv.mut.Lock() - defer srv.mut.Unlock() - if srv.schedulers != nil && srv.schedulers.running { - return ErrSchedulersRunning - } - mlog.Debug("Initialising schedulers.") - - schedulers := &Schedulers{ - stop: make(chan bool), - stopped: make(chan bool), - configChanged: make(chan *model.Config), - clusterLeaderChanged: make(chan bool, 1), - jobs: srv, - isLeader: true, - schedulers: make(map[string]model.Scheduler), - nextRunTimes: make(map[string]*time.Time), - } - - srv.schedulers = schedulers - - return nil -} - func (schedulers *Schedulers) AddScheduler(name string, scheduler model.Scheduler) { schedulers.schedulers[name] = scheduler } @@ -63,6 +39,8 @@ func (schedulers *Schedulers) AddScheduler(name string, scheduler model.Schedule // Start starts the schedulers. This call is not safe for concurrent use. // Synchronization should be implemented by the caller. func (schedulers *Schedulers) Start() { + schedulers.stop = make(chan bool) + schedulers.stopped = make(chan bool) schedulers.listenerId = schedulers.jobs.ConfigService.AddConfigListener(schedulers.handleConfigChange) go func() { diff --git a/jobs/schedulers_test.go b/jobs/schedulers_test.go index 12f68a81b3..8421d4f0f4 100644 --- a/jobs/schedulers_test.go +++ b/jobs/schedulers_test.go @@ -61,7 +61,7 @@ func TestScheduler(t *testing.T) { }, } - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.RegisterJobType(model.JobTypeDataRetention, nil, new(MockScheduler)) jobServer.RegisterJobType(model.JobTypeMessageExport, nil, new(MockScheduler)) @@ -77,7 +77,7 @@ func TestScheduler(t *testing.T) { }) t.Run("ClusterLeaderChanged", func(t *testing.T) { - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.StartSchedulers() time.Sleep(time.Second) jobServer.HandleClusterLeaderChange(false) @@ -89,7 +89,7 @@ func TestScheduler(t *testing.T) { }) t.Run("ClusterLeaderChangedBeforeStart", func(t *testing.T) { - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.HandleClusterLeaderChange(false) jobServer.StartSchedulers() time.Sleep(time.Second) @@ -100,7 +100,7 @@ func TestScheduler(t *testing.T) { }) t.Run("DoubleClusterLeaderChangedBeforeStart", func(t *testing.T) { - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.HandleClusterLeaderChange(false) jobServer.HandleClusterLeaderChange(true) jobServer.StartSchedulers() @@ -112,7 +112,7 @@ func TestScheduler(t *testing.T) { }) t.Run("ConfigChanged", func(t *testing.T) { - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.StartSchedulers() time.Sleep(time.Second) jobServer.HandleClusterLeaderChange(false) @@ -125,7 +125,7 @@ func TestScheduler(t *testing.T) { }) t.Run("ConfigChangedDeadlock", func(t *testing.T) { - jobServer.InitSchedulers() + jobServer.initSchedulers() jobServer.StartSchedulers() time.Sleep(time.Second) diff --git a/jobs/server.go b/jobs/server.go index b7e7dbe341..9c5870027c 100644 --- a/jobs/server.go +++ b/jobs/server.go @@ -5,6 +5,7 @@ package jobs import ( "sync" + "time" "github.com/mattermost/mattermost-server/v6/einterfaces" "github.com/mattermost/mattermost-server/v6/model" @@ -24,11 +25,33 @@ type JobServer struct { } func NewJobServer(configService configservice.ConfigService, store store.Store, metrics einterfaces.MetricsInterface) *JobServer { - return &JobServer{ + srv := &JobServer{ ConfigService: configService, Store: store, metrics: metrics, } + srv.initWorkers() + srv.initSchedulers() + return srv +} + +func (srv *JobServer) initWorkers() { + workers := NewWorkers(srv.ConfigService) + workers.Watcher = srv.MakeWatcher(workers, DefaultWatcherPollingInterval) + srv.workers = workers +} + +func (srv *JobServer) initSchedulers() { + schedulers := &Schedulers{ + configChanged: make(chan *model.Config), + clusterLeaderChanged: make(chan bool, 1), + jobs: srv, + isLeader: true, + schedulers: make(map[string]model.Scheduler), + nextRunTimes: make(map[string]*time.Time), + } + + srv.schedulers = schedulers } func (srv *JobServer) Config() *model.Config { diff --git a/jobs/server_test.go b/jobs/server_test.go index 391fc48d9b..9c550c54d9 100644 --- a/jobs/server_test.go +++ b/jobs/server_test.go @@ -5,44 +5,11 @@ package jobs import ( "testing" + "time" "github.com/stretchr/testify/require" ) -func TestInitWorkers(t *testing.T) { - t.Run("initialize", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - }) - - t.Run("re-initialize", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - err = jobServer.InitWorkers() - require.NoError(t, err) - }) - - t.Run("re-initialize already running", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - - err = jobServer.StartWorkers() - require.NoError(t, err) - - err = jobServer.InitWorkers() - require.Equal(t, ErrWorkersRunning, err) - - err = jobServer.StopWorkers() - require.NoError(t, err) - - err = jobServer.InitWorkers() - require.NoError(t, err) - }) -} - func TestStartWorkers(t *testing.T) { t.Run("uninitialized", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) @@ -52,22 +19,24 @@ func TestStartWorkers(t *testing.T) { t.Run("already running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - err = jobServer.StartWorkers() + jobServer.initWorkers() + err := jobServer.StartWorkers() require.NoError(t, err) err = jobServer.StartWorkers() require.Equal(t, ErrWorkersRunning, err) + // Parking the go routing to let the worker watcher start + time.Sleep(1 * time.Millisecond) err = jobServer.StopWorkers() require.NoError(t, err) }) t.Run("not running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - err = jobServer.StartWorkers() + jobServer.initWorkers() + err := jobServer.StartWorkers() require.NoError(t, err) + // Parking the go routing to let the worker watcher start + time.Sleep(1 * time.Millisecond) err = jobServer.StopWorkers() require.NoError(t, err) }) @@ -82,57 +51,23 @@ func TestStopWorkers(t *testing.T) { t.Run("not running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - err = jobServer.StopWorkers() + jobServer.initWorkers() + err := jobServer.StopWorkers() require.Equal(t, ErrWorkersNotRunning, err) }) t.Run("running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitWorkers() - require.NoError(t, err) - err = jobServer.StartWorkers() + jobServer.initWorkers() + err := jobServer.StartWorkers() require.NoError(t, err) + // Parking the go routing to let the worker watcher start + time.Sleep(1 * time.Millisecond) err = jobServer.StopWorkers() require.NoError(t, err) }) } -func TestInitSchedulers(t *testing.T) { - t.Run("initialize", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - }) - - t.Run("re-initialize", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - err = jobServer.InitSchedulers() - require.NoError(t, err) - }) - - t.Run("re-initialize already running", func(t *testing.T) { - jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - - err = jobServer.StartSchedulers() - require.NoError(t, err) - - err = jobServer.InitSchedulers() - require.Equal(t, ErrSchedulersRunning, err) - - err = jobServer.StopSchedulers() - require.NoError(t, err) - - err = jobServer.InitSchedulers() - require.NoError(t, err) - }) -} - func TestStartSchedulers(t *testing.T) { t.Run("uninitialized", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) @@ -142,9 +77,8 @@ func TestStartSchedulers(t *testing.T) { t.Run("initialized", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - err = jobServer.StartSchedulers() + jobServer.initSchedulers() + err := jobServer.StartSchedulers() require.NoError(t, err) err = jobServer.StopSchedulers() @@ -153,9 +87,8 @@ func TestStartSchedulers(t *testing.T) { t.Run("already running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - err = jobServer.StartSchedulers() + jobServer.initSchedulers() + err := jobServer.StartSchedulers() require.NoError(t, err) err = jobServer.StartSchedulers() require.Equal(t, ErrSchedulersRunning, err) @@ -174,17 +107,15 @@ func TestStopSchedulers(t *testing.T) { t.Run("not running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - err = jobServer.StopSchedulers() + jobServer.initSchedulers() + err := jobServer.StopSchedulers() require.Equal(t, ErrSchedulersNotRunning, err) }) t.Run("running", func(t *testing.T) { jobServer, _, _ := makeJobServer(t) - err := jobServer.InitSchedulers() - require.NoError(t, err) - err = jobServer.StartSchedulers() + jobServer.initSchedulers() + err := jobServer.StartSchedulers() require.NoError(t, err) err = jobServer.StopSchedulers() require.NoError(t, err) diff --git a/jobs/workers.go b/jobs/workers.go index aa6251ec72..009b040f8d 100644 --- a/jobs/workers.go +++ b/jobs/workers.go @@ -27,23 +27,6 @@ var ( ErrWorkersUninitialized = errors.New("job workers are not initialized") ) -func (srv *JobServer) InitWorkers() error { - srv.mut.Lock() - defer srv.mut.Unlock() - - if srv.workers != nil && srv.workers.running { - return ErrWorkersRunning - } - - workers := NewWorkers(srv.ConfigService) - - workers.Watcher = srv.MakeWatcher(workers, DefaultWatcherPollingInterval) - - srv.workers = workers - - return nil -} - func NewWorkers(configService configservice.ConfigService) *Workers { return &Workers{ ConfigService: configService,