diff --git a/api4/license.go b/api4/license.go index d81552acd1..89d23d658c 100644 --- a/api4/license.go +++ b/api4/license.go @@ -107,11 +107,6 @@ func addLicense(c *Context, w http.ResponseWriter, r *http.Request) { return } - if *c.App.Config().JobSettings.RunJobs { - c.App.Srv().Jobs.Workers = c.App.Srv().Jobs.InitWorkers() - c.App.Srv().Jobs.StartWorkers() - } - auditRec.Success() c.LogAudit("success") diff --git a/api4/license_local.go b/api4/license_local.go index e5cad8a389..36462f8da5 100644 --- a/api4/license_local.go +++ b/api4/license_local.go @@ -67,11 +67,6 @@ func localAddLicense(c *Context, w http.ResponseWriter, r *http.Request) { return } - if *c.App.Config().JobSettings.RunJobs { - c.App.Srv().Jobs.Workers = c.App.Srv().Jobs.InitWorkers() - c.App.Srv().Jobs.StartWorkers() - } - auditRec.Success() c.LogAudit("success") diff --git a/app/app.go b/app/app.go index 62ef79b63b..6d65ef9245 100644 --- a/app/app.go +++ b/app/app.go @@ -91,13 +91,13 @@ func (a *App) InitServer() { a.srv.ShutDownPlugins() } }) - if a.Srv().runjobs { + if a.Srv().runEssentialJobs { a.Srv().Go(func() { runLicenseExpirationCheckJob(a) runCheckWarnMetricStatusJob(a) }) + a.srv.runJobs() } - a.srv.RunJobs() }) } @@ -140,8 +140,8 @@ func (a *App) initJobs() { a.srv.Jobs.Cloud = jobsCloudInterface(a.srv) } - a.srv.Jobs.Workers = a.srv.Jobs.InitWorkers() - a.srv.Jobs.Schedulers = a.srv.Jobs.InitSchedulers() + a.srv.Jobs.InitWorkers() + a.srv.Jobs.InitSchedulers() } func (a *App) TelemetryId() string { diff --git a/app/license.go b/app/license.go index 7a5cd36d50..3619b201d8 100644 --- a/app/license.go +++ b/app/license.go @@ -13,6 +13,7 @@ import ( "github.com/dgrijalva/jwt-go" "github.com/pkg/errors" + "github.com/mattermost/mattermost-server/v5/jobs" "github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/shared/mlog" "github.com/mattermost/mattermost-server/v5/utils" @@ -124,14 +125,23 @@ func (s *Server) SaveLicense(licenseBytes []byte) (*model.License, *model.AppErr s.ReloadConfig() s.InvalidateAllCaches() - // start job server if necessary - this handles the edge case where a license file is uploaded, but the job server + // 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 && s.Jobs.Workers != nil { - s.Jobs.StartWorkers() + 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 && s.Jobs.Schedulers != nil { - s.Jobs.StartSchedulers() + 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/options.go b/app/options.go index 54c111214e..932827ceff 100644 --- a/app/options.go +++ b/app/options.go @@ -64,8 +64,8 @@ func ConfigStore(configStore *config.Store) Option { } } -func RunJobs(s *Server) error { - s.runjobs = true +func RunEssentialJobs(s *Server) error { + s.runEssentialJobs = true return nil } diff --git a/app/server.go b/app/server.go index 4ead6df570..7fcbceb13a 100644 --- a/app/server.go +++ b/app/server.go @@ -118,8 +118,8 @@ type Server struct { PushNotificationsHub PushNotificationsHub pushNotificationClient *http.Client // TODO: move this to it's own package - runjobs bool - Jobs *jobs.JobServer + runEssentialJobs bool + Jobs *jobs.JobServer clusterLeaderListeners sync.Map @@ -447,8 +447,8 @@ func NewServer(options ...Option) (*Server, error) { s.clusterLeaderListenerId = s.AddClusterLeaderChangedListener(func() { mlog.Info("Cluster leader changed. Determining if job schedulers should be running:", mlog.Bool("isLeader", s.IsLeader())) - if s.Jobs != nil && s.Jobs.Schedulers != nil { - s.Jobs.Schedulers.HandleClusterLeaderChange(s.IsLeader()) + if s.Jobs != nil { + s.Jobs.HandleClusterLeaderChange(s.IsLeader()) } s.setupFeatureFlags() }) @@ -661,44 +661,46 @@ func maxInt(a, b int) int { return b } -func (s *Server) RunJobs() { - if s.runjobs { - s.Go(func() { - runSecurityJob(s) - }) - s.Go(func() { - firstRun, err := s.getFirstServerRunTimestamp() - if err != nil { - mlog.Warn("Fetching time of first server run failed. Setting to 'now'.") - s.ensureFirstServerRunTimestamp() - firstRun = utils.MillisFromTime(time.Now()) - } - s.telemetryService.RunTelemetryJob(firstRun) - }) - s.Go(func() { - runSessionCleanupJob(s) - }) - s.Go(func() { - runTokenCleanupJob(s) - }) - s.Go(func() { - runCommandWebhookCleanupJob(s) - }) +func (s *Server) runJobs() { + s.Go(func() { + runSecurityJob(s) + }) + s.Go(func() { + firstRun, err := s.getFirstServerRunTimestamp() + if err != nil { + mlog.Warn("Fetching time of first server run failed. Setting to 'now'.") + s.ensureFirstServerRunTimestamp() + firstRun = utils.MillisFromTime(time.Now()) + } + s.telemetryService.RunTelemetryJob(firstRun) + }) + s.Go(func() { + runSessionCleanupJob(s) + }) + s.Go(func() { + runTokenCleanupJob(s) + }) + s.Go(func() { + runCommandWebhookCleanupJob(s) + }) - if complianceI := s.Compliance; complianceI != nil { - complianceI.StartComplianceDailyJob() - } + if complianceI := s.Compliance; complianceI != nil { + complianceI.StartComplianceDailyJob() + } - if *s.Config().JobSettings.RunJobs && s.Jobs != nil { - s.Jobs.StartWorkers() + if *s.Config().JobSettings.RunJobs && s.Jobs != nil { + if err := s.Jobs.StartWorkers(); err != nil { + mlog.Error("Failed to start job server workers", mlog.Err(err)) } - if *s.Config().JobSettings.RunScheduler && s.Jobs != nil { - s.Jobs.StartSchedulers() + } + if *s.Config().JobSettings.RunScheduler && s.Jobs != nil { + if err := s.Jobs.StartSchedulers(); err != nil { + mlog.Error("Failed to start job server schedulers", mlog.Err(err)) } + } - if *s.Config().ServiceSettings.EnableAWSMetering { - runReportToAWSMeterJob(s) - } + if *s.Config().ServiceSettings.EnableAWSMetering { + runReportToAWSMeterJob(s) } } @@ -895,9 +897,16 @@ func (s *Server) Shutdown() { s.StopMetricsServer() // This must be done after the cluster is stopped. - if s.Jobs != nil && s.runjobs { - s.Jobs.StopWorkers() - s.Jobs.StopSchedulers() + if s.Jobs != nil { + // For simplicity we don't check if workers and schedulers are active + // before stopping them as both calls essentially become no-ops + // if nothing is running. + if err = s.Jobs.StopWorkers(); err != nil && !errors.Is(err, jobs.ErrWorkersNotRunning) { + mlog.Warn("Failed to stop job server workers", mlog.Err(err)) + } + if err = s.Jobs.StopSchedulers(); err != nil && !errors.Is(err, jobs.ErrSchedulersNotRunning) { + mlog.Warn("Failed to stop job server schedulers", mlog.Err(err)) + } } if s.Store != nil { diff --git a/cmd/mattermost/commands/server.go b/cmd/mattermost/commands/server.go index 2c05bb5f83..fb102a03e5 100644 --- a/cmd/mattermost/commands/server.go +++ b/cmd/mattermost/commands/server.go @@ -69,7 +69,7 @@ func runServer(configStore *config.Store, usedPlatform bool, interruptChan chan options := []app.Option{ app.ConfigStore(configStore), - app.RunJobs, + app.RunEssentialJobs, app.JoinCluster, app.StartSearchEngine, app.StartMetrics, diff --git a/jobs/schedulers.go b/jobs/schedulers.go index 243937c225..4c829ac9a2 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -4,8 +4,8 @@ package jobs import ( + "errors" "fmt" - "sync" "time" "github.com/mattermost/mattermost-server/v5/model" @@ -18,15 +18,26 @@ type Schedulers struct { configChanged chan *model.Config clusterLeaderChanged chan bool listenerId string - startOnce sync.Once jobs *JobServer isLeader bool + running bool schedulers []model.Scheduler nextRunTimes []*time.Time } -func (srv *JobServer) InitSchedulers() *Schedulers { +var ( + ErrSchedulersNotRunning = errors.New("job schedulers are not running") + ErrSchedulersRunning = errors.New("job schedulers are running") + 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{ @@ -87,89 +98,94 @@ func (srv *JobServer) InitSchedulers() *Schedulers { } schedulers.nextRunTimes = make([]*time.Time, len(schedulers.schedulers)) - return schedulers + srv.schedulers = schedulers + + return nil } -func (schedulers *Schedulers) Start() *Schedulers { +// Start starts the schedulers. This call is not safe for concurrent use. +// Synchronization should be implemented by the caller. +func (schedulers *Schedulers) Start() { schedulers.listenerId = schedulers.jobs.ConfigService.AddConfigListener(schedulers.handleConfigChange) go func() { - schedulers.startOnce.Do(func() { - mlog.Info("Starting schedulers.") + mlog.Info("Starting schedulers.") - defer func() { - mlog.Info("Schedulers stopped.") - close(schedulers.stopped) - }() + defer func() { + mlog.Info("Schedulers stopped.") + close(schedulers.stopped) + }() - now := time.Now() - for idx, scheduler := range schedulers.schedulers { - if !scheduler.Enabled(schedulers.jobs.Config()) { - schedulers.nextRunTimes[idx] = nil - } else { - schedulers.setNextRunTime(schedulers.jobs.Config(), idx, now, false) - } + now := time.Now() + for idx, scheduler := range schedulers.schedulers { + if !scheduler.Enabled(schedulers.jobs.Config()) { + schedulers.nextRunTimes[idx] = nil + } else { + schedulers.setNextRunTime(schedulers.jobs.Config(), idx, now, false) } + } - for { - timer := time.NewTimer(1 * time.Minute) - select { - case <-schedulers.stop: - mlog.Debug("Schedulers received stop signal.") - timer.Stop() - return - case now = <-timer.C: - cfg := schedulers.jobs.Config() + for { + timer := time.NewTimer(1 * time.Minute) + select { + case <-schedulers.stop: + mlog.Debug("Schedulers received stop signal.") + timer.Stop() + return + case now = <-timer.C: + cfg := schedulers.jobs.Config() - for idx, nextTime := range schedulers.nextRunTimes { - if nextTime == nil { + for idx, nextTime := range schedulers.nextRunTimes { + if nextTime == nil { + continue + } + + if time.Now().After(*nextTime) { + scheduler := schedulers.schedulers[idx] + if scheduler == nil || !schedulers.isLeader || !scheduler.Enabled(cfg) { continue } - - if time.Now().After(*nextTime) { - scheduler := schedulers.schedulers[idx] - if scheduler != nil { - 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 { - schedulers.setNextRunTime(cfg, idx, now, true) - } - } - } - } - } - case newCfg := <-schedulers.configChanged: - for idx, scheduler := range schedulers.schedulers { - if !schedulers.isLeader || !scheduler.Enabled(newCfg) { - schedulers.nextRunTimes[idx] = nil - } else { - schedulers.setNextRunTime(newCfg, idx, now, false) - } - } - case isLeader := <-schedulers.clusterLeaderChanged: - for idx := range schedulers.schedulers { - schedulers.isLeader = isLeader - if !isLeader { - schedulers.nextRunTimes[idx] = nil - } else { - schedulers.setNextRunTime(schedulers.jobs.Config(), idx, now, false) + if _, err := schedulers.scheduleJob(cfg, scheduler); err != nil { + mlog.Error("Failed to schedule job", mlog.String("scheduler", scheduler.Name()), mlog.Err(err)) + continue } + schedulers.setNextRunTime(cfg, idx, now, true) + } + } + case newCfg := <-schedulers.configChanged: + for idx, scheduler := range schedulers.schedulers { + if !schedulers.isLeader || !scheduler.Enabled(newCfg) { + schedulers.nextRunTimes[idx] = nil + } else { + schedulers.setNextRunTime(newCfg, idx, now, false) + } + } + case isLeader := <-schedulers.clusterLeaderChanged: + for idx := range schedulers.schedulers { + schedulers.isLeader = isLeader + if !isLeader { + schedulers.nextRunTimes[idx] = nil + } else { + schedulers.setNextRunTime(schedulers.jobs.Config(), idx, now, false) } } - timer.Stop() } - }) + timer.Stop() + } }() - return schedulers + schedulers.running = true } -func (schedulers *Schedulers) Stop() *Schedulers { +// Stop stops the schedulers. This call is not safe for concurrent use. +// Synchronization should be implemented by the caller. +func (schedulers *Schedulers) Stop() { mlog.Info("Stopping schedulers.") close(schedulers.stop) <-schedulers.stopped - return schedulers + schedulers.jobs.ConfigService.RemoveConfigListener(schedulers.listenerId) + schedulers.listenerId = "" + schedulers.running = false } func (schedulers *Schedulers) setNextRunTime(cfg *model.Config, idx int, now time.Time, pendingJobs bool) { diff --git a/jobs/schedulers_test.go b/jobs/schedulers_test.go index 13d51d7e8b..b457bc663b 100644 --- a/jobs/schedulers_test.go +++ b/jobs/schedulers_test.go @@ -78,38 +78,38 @@ func TestScheduler(t *testing.T) { jobServer.MessageExportJob = exportInterface t.Run("Base", func(t *testing.T) { - schedulers := jobServer.InitSchedulers() - schedulers.Start() + jobServer.InitSchedulers() + jobServer.StartSchedulers() time.Sleep(time.Second) - schedulers.Stop() + jobServer.StopSchedulers() // They should be all on here - for _, element := range schedulers.nextRunTimes { + for _, element := range jobServer.schedulers.nextRunTimes { assert.NotNil(t, element) } }) t.Run("ClusterLeaderChanged", func(t *testing.T) { - schedulers := jobServer.InitSchedulers() - schedulers.Start() + jobServer.InitSchedulers() + jobServer.StartSchedulers() time.Sleep(time.Second) - schedulers.HandleClusterLeaderChange(false) - schedulers.Stop() + jobServer.HandleClusterLeaderChange(false) + jobServer.StopSchedulers() // They should be turned off - for _, element := range schedulers.nextRunTimes { + for _, element := range jobServer.schedulers.nextRunTimes { assert.Nil(t, element) } }) t.Run("ConfigChanged", func(t *testing.T) { - schedulers := jobServer.InitSchedulers() - schedulers.Start() + jobServer.InitSchedulers() + jobServer.StartSchedulers() time.Sleep(time.Second) - schedulers.HandleClusterLeaderChange(false) + jobServer.HandleClusterLeaderChange(false) // After running a config change, they should stay off - schedulers.handleConfigChange(nil, nil) - schedulers.Stop() - for _, element := range schedulers.nextRunTimes { + jobServer.schedulers.handleConfigChange(nil, nil) + jobServer.StopSchedulers() + for _, element := range jobServer.schedulers.nextRunTimes { assert.Nil(t, element) } }) diff --git a/jobs/server.go b/jobs/server.go index 76f6a91af3..7e3a352421 100644 --- a/jobs/server.go +++ b/jobs/server.go @@ -4,6 +4,8 @@ package jobs import ( + "sync" + "github.com/mattermost/mattermost-server/v5/einterfaces" ejobs "github.com/mattermost/mattermost-server/v5/einterfaces/jobs" tjobs "github.com/mattermost/mattermost-server/v5/jobs/interfaces" @@ -16,8 +18,6 @@ type JobServer struct { ConfigService configservice.ConfigService Store store.Store metrics einterfaces.MetricsInterface - Workers *Workers - Schedulers *Schedulers DataRetentionJob ejobs.DataRetentionJobInterface MessageExportJob ejobs.MessageExportJobInterface @@ -35,6 +35,11 @@ type JobServer struct { ExportProcess tjobs.ExportProcessInterface ExportDelete tjobs.ExportDeleteInterface Cloud ejobs.CloudJobInterface + + // mut is used to protect the following fields from concurrent access. + mut sync.Mutex + workers *Workers + schedulers *Schedulers } func NewJobServer(configService configservice.ConfigService, store store.Store, metrics einterfaces.MetricsInterface) *JobServer { @@ -49,22 +54,58 @@ func (srv *JobServer) Config() *model.Config { return srv.ConfigService.Config() } -func (srv *JobServer) StartWorkers() { - srv.Workers = srv.Workers.Start() +func (srv *JobServer) StartWorkers() error { + srv.mut.Lock() + defer srv.mut.Unlock() + if srv.workers == nil { + return ErrWorkersUninitialized + } else if srv.workers.running { + return ErrWorkersRunning + } + srv.workers.Start() + return nil } -func (srv *JobServer) StartSchedulers() { - srv.Schedulers = srv.Schedulers.Start() +func (srv *JobServer) StartSchedulers() error { + srv.mut.Lock() + defer srv.mut.Unlock() + if srv.schedulers == nil { + return ErrSchedulersUninitialized + } else if srv.schedulers.running { + return ErrSchedulersRunning + } + srv.schedulers.Start() + return nil } -func (srv *JobServer) StopWorkers() { - if srv.Workers != nil { - srv.Workers.Stop() - } -} - -func (srv *JobServer) StopSchedulers() { - if srv.Schedulers != nil { - srv.Schedulers.Stop() +func (srv *JobServer) StopWorkers() error { + srv.mut.Lock() + defer srv.mut.Unlock() + if srv.workers == nil { + return ErrWorkersUninitialized + } else if !srv.workers.running { + return ErrWorkersNotRunning + } + srv.workers.Stop() + return nil +} + +func (srv *JobServer) StopSchedulers() error { + srv.mut.Lock() + defer srv.mut.Unlock() + if srv.schedulers == nil { + return ErrSchedulersUninitialized + } else if !srv.schedulers.running { + return ErrSchedulersNotRunning + } + srv.schedulers.Stop() + return nil +} + +func (srv *JobServer) HandleClusterLeaderChange(isLeader bool) { + srv.mut.Lock() + defer srv.mut.Unlock() + if srv.schedulers != nil { + srv.schedulers.HandleClusterLeaderChange(isLeader) } } diff --git a/jobs/server_test.go b/jobs/server_test.go new file mode 100644 index 0000000000..391fc48d9b --- /dev/null +++ b/jobs/server_test.go @@ -0,0 +1,192 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. + +package jobs + +import ( + "testing" + + "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) + err := jobServer.StartWorkers() + require.Equal(t, ErrWorkersUninitialized, err) + }) + + t.Run("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.StartWorkers() + require.Equal(t, ErrWorkersRunning, err) + 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() + require.NoError(t, err) + err = jobServer.StopWorkers() + require.NoError(t, err) + }) +} + +func TestStopWorkers(t *testing.T) { + t.Run("uninitialized", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + err := jobServer.StopWorkers() + require.Equal(t, ErrWorkersUninitialized, err) + }) + + t.Run("not running", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + err := jobServer.InitWorkers() + require.NoError(t, err) + 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() + require.NoError(t, err) + 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) + err := jobServer.StartSchedulers() + require.Equal(t, ErrSchedulersUninitialized, err) + }) + + t.Run("initialized", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + err := jobServer.InitSchedulers() + require.NoError(t, err) + err = jobServer.StartSchedulers() + require.NoError(t, err) + + err = jobServer.StopSchedulers() + require.NoError(t, err) + }) + + t.Run("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.StartSchedulers() + require.Equal(t, ErrSchedulersRunning, err) + + err = jobServer.StopSchedulers() + require.NoError(t, err) + }) +} + +func TestStopSchedulers(t *testing.T) { + t.Run("uninitialized", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + err := jobServer.StopSchedulers() + require.Equal(t, ErrSchedulersUninitialized, err) + }) + + t.Run("not running", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + err := jobServer.InitSchedulers() + require.NoError(t, err) + 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() + require.NoError(t, err) + err = jobServer.StopSchedulers() + require.NoError(t, err) + }) +} diff --git a/jobs/workers.go b/jobs/workers.go index a68ef15c73..a7d1bfb995 100644 --- a/jobs/workers.go +++ b/jobs/workers.go @@ -4,7 +4,7 @@ package jobs import ( - "sync" + "errors" "github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/services/configservice" @@ -12,7 +12,6 @@ import ( ) type Workers struct { - startOnce sync.Once ConfigService configservice.ConfigService Watcher *Watcher @@ -34,9 +33,23 @@ type Workers struct { Cloud model.Worker listenerId string + running bool } -func (srv *JobServer) InitWorkers() *Workers { +var ( + ErrWorkersNotRunning = errors.New("job workers are not running") + ErrWorkersRunning = errors.New("job workers are running") + 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 := &Workers{ ConfigService: srv.ConfigService, } @@ -106,83 +119,83 @@ func (srv *JobServer) InitWorkers() *Workers { workers.Cloud = cloudInterface.MakeWorker() } - return workers + srv.workers = workers + + return nil } -func (workers *Workers) Start() *Workers { +// Start starts the workers. This call is not safe for concurrent use. +// Synchronization should be implemented by the caller. +func (workers *Workers) Start() { mlog.Info("Starting workers") + if workers.DataRetention != nil && (*workers.ConfigService.Config().DataRetentionSettings.EnableMessageDeletion || *workers.ConfigService.Config().DataRetentionSettings.EnableFileDeletion) { + go workers.DataRetention.Run() + } - workers.startOnce.Do(func() { - if workers.DataRetention != nil && (*workers.ConfigService.Config().DataRetentionSettings.EnableMessageDeletion || *workers.ConfigService.Config().DataRetentionSettings.EnableFileDeletion) { - go workers.DataRetention.Run() - } + if workers.MessageExport != nil && *workers.ConfigService.Config().MessageExportSettings.EnableExport { + go workers.MessageExport.Run() + } - if workers.MessageExport != nil && *workers.ConfigService.Config().MessageExportSettings.EnableExport { - go workers.MessageExport.Run() - } + if workers.ElasticsearchIndexing != nil && *workers.ConfigService.Config().ElasticsearchSettings.EnableIndexing { + go workers.ElasticsearchIndexing.Run() + } - if workers.ElasticsearchIndexing != nil && *workers.ConfigService.Config().ElasticsearchSettings.EnableIndexing { - go workers.ElasticsearchIndexing.Run() - } + if workers.ElasticsearchAggregation != nil && *workers.ConfigService.Config().ElasticsearchSettings.EnableIndexing { + go workers.ElasticsearchAggregation.Run() + } - if workers.ElasticsearchAggregation != nil && *workers.ConfigService.Config().ElasticsearchSettings.EnableIndexing { - go workers.ElasticsearchAggregation.Run() - } + if workers.LdapSync != nil && *workers.ConfigService.Config().LdapSettings.EnableSync { + go workers.LdapSync.Run() + } - if workers.LdapSync != nil && *workers.ConfigService.Config().LdapSettings.EnableSync { - go workers.LdapSync.Run() - } + if workers.Migrations != nil { + go workers.Migrations.Run() + } - if workers.Migrations != nil { - go workers.Migrations.Run() - } + if workers.Plugins != nil { + go workers.Plugins.Run() + } - if workers.Plugins != nil { - go workers.Plugins.Run() - } + if workers.BleveIndexing != nil && *workers.ConfigService.Config().BleveSettings.EnableIndexing && *workers.ConfigService.Config().BleveSettings.IndexDir != "" { + go workers.BleveIndexing.Run() + } - if workers.BleveIndexing != nil && *workers.ConfigService.Config().BleveSettings.EnableIndexing && *workers.ConfigService.Config().BleveSettings.IndexDir != "" { - go workers.BleveIndexing.Run() - } + if workers.ExpiryNotify != nil { + go workers.ExpiryNotify.Run() + } - if workers.ExpiryNotify != nil { - go workers.ExpiryNotify.Run() - } + if workers.ActiveUsers != nil { + go workers.ActiveUsers.Run() + } - if workers.ActiveUsers != nil { - go workers.ActiveUsers.Run() - } + if workers.ProductNotices != nil { + go workers.ProductNotices.Run() + } - if workers.ProductNotices != nil { - go workers.ProductNotices.Run() - } + if workers.ImportProcess != nil { + go workers.ImportProcess.Run() + } - if workers.ImportProcess != nil { - go workers.ImportProcess.Run() - } + if workers.ImportDelete != nil { + go workers.ImportDelete.Run() + } - if workers.ImportDelete != nil { - go workers.ImportDelete.Run() - } + if workers.ExportProcess != nil { + go workers.ExportProcess.Run() + } - if workers.ExportProcess != nil { - go workers.ExportProcess.Run() - } + if workers.ExportDelete != nil { + go workers.ExportDelete.Run() + } - if workers.ExportDelete != nil { - go workers.ExportDelete.Run() - } + if workers.Cloud != nil { + go workers.Cloud.Run() + } - if workers.Cloud != nil { - go workers.Cloud.Run() - } - - go workers.Watcher.Start() - }) + go workers.Watcher.Start() workers.listenerId = workers.ConfigService.AddConfigListener(workers.handleConfigChange) - - return workers + workers.running = true } func (workers *Workers) handleConfigChange(oldConfig *model.Config, newConfig *model.Config) { @@ -237,7 +250,9 @@ func (workers *Workers) handleConfigChange(oldConfig *model.Config, newConfig *m } } -func (workers *Workers) Stop() *Workers { +// Stop stops the workers. This call is not safe for concurrent use. +// Synchronization should be implemented by the caller. +func (workers *Workers) Stop() { workers.ConfigService.RemoveConfigListener(workers.listenerId) workers.Watcher.Stop() @@ -306,7 +321,7 @@ func (workers *Workers) Stop() *Workers { workers.Cloud.Stop() } - mlog.Info("Stopped workers") + workers.running = false - return workers + mlog.Info("Stopped workers") }