коммит произвёл
GitHub
родитель
a2a78577e9
Коммит
4deedbbfd3
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
Ссылка в новой задаче
Block a user