[MM-35424] Fix job schedulers server from missing leader change event (#17602)
* Fix cluster leader change message potentially getting lost * Start jobs and schedulers earlier in server initialization * Fix possible deadlock * Update jobs/schedulers.go Co-authored-by: Ibrahim Serdar Acikgoz <serdaracikgoz86@gmail.com> Co-authored-by: Ibrahim Serdar Acikgoz <serdaracikgoz86@gmail.com>
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
4c4d889739
Коммит
7bb23323d1
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user