diff --git a/app/server.go b/app/server.go index fd3425b8e6..d3e44db352 100644 --- a/app/server.go +++ b/app/server.go @@ -662,6 +662,16 @@ func NewServer(options ...Option) (*Server, error) { } } + if s.runEssentialJobs { + s.Go(func() { + s.runLicenseExpirationCheckJob() + runCheckAdminSupportStatusJob(fakeApp, c) + runCheckWarnMetricStatusJob(fakeApp, c) + runDNDStatusExpireJob(fakeApp) + }) + s.runJobs() + } + s.doAppMigrations() s.initPostMetadata() @@ -674,15 +684,6 @@ func NewServer(options ...Option) (*Server, error) { s.ShutDownPlugins() } }) - if s.runEssentialJobs { - s.Go(func() { - s.runLicenseExpirationCheckJob() - runCheckAdminSupportStatusJob(fakeApp, c) - runCheckWarnMetricStatusJob(fakeApp, c) - runDNDStatusExpireJob(fakeApp) - }) - s.runJobs() - } return s, nil } diff --git a/jobs/schedulers.go b/jobs/schedulers.go index 7a581d0f18..391c1176ee 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -44,7 +44,7 @@ func (srv *JobServer) InitSchedulers() error { stop: make(chan bool), stopped: make(chan bool), configChanged: make(chan *model.Config), - clusterLeaderChanged: make(chan bool), + clusterLeaderChanged: make(chan bool, 1), jobs: srv, isLeader: true, } @@ -230,15 +230,28 @@ func (schedulers *Schedulers) scheduleJob(cfg *model.Config, scheduler model.Sch return scheduler.ScheduleJob(cfg, pendingJobs, lastSuccessfulJob) } -func (schedulers *Schedulers) handleConfigChange(oldConfig, newConfig *model.Config) { +func (schedulers *Schedulers) handleConfigChange(_, newConfig *model.Config) { mlog.Debug("Schedulers received config change.") - schedulers.configChanged <- newConfig + select { + case schedulers.configChanged <- newConfig: + case <-schedulers.stop: + } } -func (schedulers *Schedulers) HandleClusterLeaderChange(isLeader bool) { +func (schedulers *Schedulers) handleClusterLeaderChange(isLeader bool) { select { case schedulers.clusterLeaderChanged <- isLeader: default: - mlog.Debug("Did not send cluster leader change message to schedulers as no schedulers listening to notification channel.") + mlog.Debug("Sending cluster leader change message to schedulers failed.") + + // Drain the buffered channel to make room for the latest change. + select { + case <-schedulers.clusterLeaderChanged: + default: + } + + // Enqueue the latest change. This operation is safe due to this method + // being called under lock. + schedulers.clusterLeaderChanged <- isLeader } } diff --git a/jobs/schedulers_test.go b/jobs/schedulers_test.go index b457bc663b..6050a5fcc1 100644 --- a/jobs/schedulers_test.go +++ b/jobs/schedulers_test.go @@ -3,6 +3,7 @@ package jobs import ( + "sync" "testing" "time" @@ -101,6 +102,29 @@ func TestScheduler(t *testing.T) { } }) + t.Run("ClusterLeaderChangedBeforeStart", func(t *testing.T) { + jobServer.InitSchedulers() + jobServer.HandleClusterLeaderChange(false) + jobServer.StartSchedulers() + time.Sleep(time.Second) + for _, element := range jobServer.schedulers.nextRunTimes { + assert.Nil(t, element) + } + jobServer.StopSchedulers() + }) + + t.Run("DoubleClusterLeaderChangedBeforeStart", func(t *testing.T) { + jobServer.InitSchedulers() + jobServer.HandleClusterLeaderChange(false) + jobServer.HandleClusterLeaderChange(true) + jobServer.StartSchedulers() + time.Sleep(time.Second) + for _, element := range jobServer.schedulers.nextRunTimes { + assert.NotNil(t, element) + } + jobServer.StopSchedulers() + }) + t.Run("ConfigChanged", func(t *testing.T) { jobServer.InitSchedulers() jobServer.StartSchedulers() @@ -113,4 +137,23 @@ func TestScheduler(t *testing.T) { assert.Nil(t, element) } }) + + t.Run("ConfigChangedDeadlock", func(t *testing.T) { + jobServer.InitSchedulers() + jobServer.StartSchedulers() + time.Sleep(time.Second) + + var wg sync.WaitGroup + wg.Add(2) + go func() { + defer wg.Done() + jobServer.StopSchedulers() + }() + go func() { + defer wg.Done() + jobServer.schedulers.handleConfigChange(nil, nil) + }() + + wg.Wait() + }) } diff --git a/jobs/server.go b/jobs/server.go index 5ccc3dc46d..74703e4397 100644 --- a/jobs/server.go +++ b/jobs/server.go @@ -107,6 +107,6 @@ func (srv *JobServer) HandleClusterLeaderChange(isLeader bool) { srv.mut.Lock() defer srv.mut.Unlock() if srv.schedulers != nil { - srv.schedulers.HandleClusterLeaderChange(isLeader) + srv.schedulers.handleClusterLeaderChange(isLeader) } }