[MM-32822] Improve JobServer state handling (#16959)

Automatic Merge
Этот коммит содержится в:
Claudio Costa
2021-03-26 10:32:16 +01:00
коммит произвёл GitHub
родитель 85bc354680
Коммит 017e14d4b5
12 изменённых файлов: 490 добавлений и 217 удалений

Просмотреть файл

@@ -107,11 +107,6 @@ func addLicense(c *Context, w http.ResponseWriter, r *http.Request) {
return return
} }
if *c.App.Config().JobSettings.RunJobs {
c.App.Srv().Jobs.Workers = c.App.Srv().Jobs.InitWorkers()
c.App.Srv().Jobs.StartWorkers()
}
auditRec.Success() auditRec.Success()
c.LogAudit("success") c.LogAudit("success")

Просмотреть файл

@@ -67,11 +67,6 @@ func localAddLicense(c *Context, w http.ResponseWriter, r *http.Request) {
return return
} }
if *c.App.Config().JobSettings.RunJobs {
c.App.Srv().Jobs.Workers = c.App.Srv().Jobs.InitWorkers()
c.App.Srv().Jobs.StartWorkers()
}
auditRec.Success() auditRec.Success()
c.LogAudit("success") c.LogAudit("success")

Просмотреть файл

@@ -91,13 +91,13 @@ func (a *App) InitServer() {
a.srv.ShutDownPlugins() a.srv.ShutDownPlugins()
} }
}) })
if a.Srv().runjobs { if a.Srv().runEssentialJobs {
a.Srv().Go(func() { a.Srv().Go(func() {
runLicenseExpirationCheckJob(a) runLicenseExpirationCheckJob(a)
runCheckWarnMetricStatusJob(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.Cloud = jobsCloudInterface(a.srv)
} }
a.srv.Jobs.Workers = a.srv.Jobs.InitWorkers() a.srv.Jobs.InitWorkers()
a.srv.Jobs.Schedulers = a.srv.Jobs.InitSchedulers() a.srv.Jobs.InitSchedulers()
} }
func (a *App) TelemetryId() string { func (a *App) TelemetryId() string {

Просмотреть файл

@@ -13,6 +13,7 @@ import (
"github.com/dgrijalva/jwt-go" "github.com/dgrijalva/jwt-go"
"github.com/pkg/errors" "github.com/pkg/errors"
"github.com/mattermost/mattermost-server/v5/jobs"
"github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/model"
"github.com/mattermost/mattermost-server/v5/shared/mlog" "github.com/mattermost/mattermost-server/v5/shared/mlog"
"github.com/mattermost/mattermost-server/v5/utils" "github.com/mattermost/mattermost-server/v5/utils"
@@ -124,14 +125,23 @@ func (s *Server) SaveLicense(licenseBytes []byte) (*model.License, *model.AppErr
s.ReloadConfig() s.ReloadConfig()
s.InvalidateAllCaches() 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 // doesn't start until the server is restarted, which prevents the 'run job now' buttons in system console from
// functioning as expected // functioning as expected
if *s.Config().JobSettings.RunJobs && s.Jobs != nil && s.Jobs.Workers != nil { if *s.Config().JobSettings.RunJobs && s.Jobs != nil {
s.Jobs.StartWorkers() 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))
} }
if *s.Config().JobSettings.RunScheduler && s.Jobs != nil && s.Jobs.Schedulers != nil {
s.Jobs.StartSchedulers()
} }
return license, nil return license, nil

Просмотреть файл

@@ -64,8 +64,8 @@ func ConfigStore(configStore *config.Store) Option {
} }
} }
func RunJobs(s *Server) error { func RunEssentialJobs(s *Server) error {
s.runjobs = true s.runEssentialJobs = true
return nil return nil
} }

Просмотреть файл

@@ -118,7 +118,7 @@ type Server struct {
PushNotificationsHub PushNotificationsHub PushNotificationsHub PushNotificationsHub
pushNotificationClient *http.Client // TODO: move this to it's own package pushNotificationClient *http.Client // TODO: move this to it's own package
runjobs bool runEssentialJobs bool
Jobs *jobs.JobServer Jobs *jobs.JobServer
clusterLeaderListeners sync.Map clusterLeaderListeners sync.Map
@@ -447,8 +447,8 @@ func NewServer(options ...Option) (*Server, error) {
s.clusterLeaderListenerId = s.AddClusterLeaderChangedListener(func() { s.clusterLeaderListenerId = s.AddClusterLeaderChangedListener(func() {
mlog.Info("Cluster leader changed. Determining if job schedulers should be running:", mlog.Bool("isLeader", s.IsLeader())) 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 { if s.Jobs != nil {
s.Jobs.Schedulers.HandleClusterLeaderChange(s.IsLeader()) s.Jobs.HandleClusterLeaderChange(s.IsLeader())
} }
s.setupFeatureFlags() s.setupFeatureFlags()
}) })
@@ -661,8 +661,7 @@ func maxInt(a, b int) int {
return b return b
} }
func (s *Server) RunJobs() { func (s *Server) runJobs() {
if s.runjobs {
s.Go(func() { s.Go(func() {
runSecurityJob(s) runSecurityJob(s)
}) })
@@ -690,17 +689,20 @@ func (s *Server) RunJobs() {
} }
if *s.Config().JobSettings.RunJobs && s.Jobs != nil { if *s.Config().JobSettings.RunJobs && s.Jobs != nil {
s.Jobs.StartWorkers() 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 { if *s.Config().JobSettings.RunScheduler && s.Jobs != nil {
s.Jobs.StartSchedulers() if err := s.Jobs.StartSchedulers(); err != nil {
mlog.Error("Failed to start job server schedulers", mlog.Err(err))
}
} }
if *s.Config().ServiceSettings.EnableAWSMetering { if *s.Config().ServiceSettings.EnableAWSMetering {
runReportToAWSMeterJob(s) runReportToAWSMeterJob(s)
} }
} }
}
// Global app options that should be applied to apps created by this server // Global app options that should be applied to apps created by this server
func (s *Server) AppOptions() []AppOption { func (s *Server) AppOptions() []AppOption {
@@ -895,9 +897,16 @@ func (s *Server) Shutdown() {
s.StopMetricsServer() s.StopMetricsServer()
// This must be done after the cluster is stopped. // This must be done after the cluster is stopped.
if s.Jobs != nil && s.runjobs { if s.Jobs != nil {
s.Jobs.StopWorkers() // For simplicity we don't check if workers and schedulers are active
s.Jobs.StopSchedulers() // 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 { if s.Store != nil {

Просмотреть файл

@@ -69,7 +69,7 @@ func runServer(configStore *config.Store, usedPlatform bool, interruptChan chan
options := []app.Option{ options := []app.Option{
app.ConfigStore(configStore), app.ConfigStore(configStore),
app.RunJobs, app.RunEssentialJobs,
app.JoinCluster, app.JoinCluster,
app.StartSearchEngine, app.StartSearchEngine,
app.StartMetrics, app.StartMetrics,

Просмотреть файл

@@ -4,8 +4,8 @@
package jobs package jobs
import ( import (
"errors"
"fmt" "fmt"
"sync"
"time" "time"
"github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/model"
@@ -18,15 +18,26 @@ type Schedulers struct {
configChanged chan *model.Config configChanged chan *model.Config
clusterLeaderChanged chan bool clusterLeaderChanged chan bool
listenerId string listenerId string
startOnce sync.Once
jobs *JobServer jobs *JobServer
isLeader bool isLeader bool
running bool
schedulers []model.Scheduler schedulers []model.Scheduler
nextRunTimes []*time.Time 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.") mlog.Debug("Initialising schedulers.")
schedulers := &Schedulers{ schedulers := &Schedulers{
@@ -87,14 +98,17 @@ func (srv *JobServer) InitSchedulers() *Schedulers {
} }
schedulers.nextRunTimes = make([]*time.Time, len(schedulers.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) schedulers.listenerId = schedulers.jobs.ConfigService.AddConfigListener(schedulers.handleConfigChange)
go func() { go func() {
schedulers.startOnce.Do(func() {
mlog.Info("Starting schedulers.") mlog.Info("Starting schedulers.")
defer func() { defer func() {
@@ -128,17 +142,16 @@ func (schedulers *Schedulers) Start() *Schedulers {
if time.Now().After(*nextTime) { if time.Now().After(*nextTime) {
scheduler := schedulers.schedulers[idx] scheduler := schedulers.schedulers[idx]
if scheduler != nil { if scheduler == nil || !schedulers.isLeader || !scheduler.Enabled(cfg) {
if schedulers.isLeader && scheduler.Enabled(cfg) { continue
}
if _, err := schedulers.scheduleJob(cfg, scheduler); err != nil { if _, err := schedulers.scheduleJob(cfg, scheduler); err != nil {
mlog.Error("Failed to schedule job", mlog.String("scheduler", scheduler.Name()), mlog.Err(err)) mlog.Error("Failed to schedule job", mlog.String("scheduler", scheduler.Name()), mlog.Err(err))
} else { continue
}
schedulers.setNextRunTime(cfg, idx, now, true) schedulers.setNextRunTime(cfg, idx, now, true)
} }
} }
}
}
}
case newCfg := <-schedulers.configChanged: case newCfg := <-schedulers.configChanged:
for idx, scheduler := range schedulers.schedulers { for idx, scheduler := range schedulers.schedulers {
if !schedulers.isLeader || !scheduler.Enabled(newCfg) { if !schedulers.isLeader || !scheduler.Enabled(newCfg) {
@@ -159,17 +172,20 @@ func (schedulers *Schedulers) Start() *Schedulers {
} }
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.") mlog.Info("Stopping schedulers.")
close(schedulers.stop) close(schedulers.stop)
<-schedulers.stopped <-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) { func (schedulers *Schedulers) setNextRunTime(cfg *model.Config, idx int, now time.Time, pendingJobs bool) {

Просмотреть файл

@@ -78,38 +78,38 @@ func TestScheduler(t *testing.T) {
jobServer.MessageExportJob = exportInterface jobServer.MessageExportJob = exportInterface
t.Run("Base", func(t *testing.T) { t.Run("Base", func(t *testing.T) {
schedulers := jobServer.InitSchedulers() jobServer.InitSchedulers()
schedulers.Start() jobServer.StartSchedulers()
time.Sleep(time.Second) time.Sleep(time.Second)
schedulers.Stop() jobServer.StopSchedulers()
// They should be all on here // They should be all on here
for _, element := range schedulers.nextRunTimes { for _, element := range jobServer.schedulers.nextRunTimes {
assert.NotNil(t, element) assert.NotNil(t, element)
} }
}) })
t.Run("ClusterLeaderChanged", func(t *testing.T) { t.Run("ClusterLeaderChanged", func(t *testing.T) {
schedulers := jobServer.InitSchedulers() jobServer.InitSchedulers()
schedulers.Start() jobServer.StartSchedulers()
time.Sleep(time.Second) time.Sleep(time.Second)
schedulers.HandleClusterLeaderChange(false) jobServer.HandleClusterLeaderChange(false)
schedulers.Stop() jobServer.StopSchedulers()
// They should be turned off // They should be turned off
for _, element := range schedulers.nextRunTimes { for _, element := range jobServer.schedulers.nextRunTimes {
assert.Nil(t, element) assert.Nil(t, element)
} }
}) })
t.Run("ConfigChanged", func(t *testing.T) { t.Run("ConfigChanged", func(t *testing.T) {
schedulers := jobServer.InitSchedulers() jobServer.InitSchedulers()
schedulers.Start() jobServer.StartSchedulers()
time.Sleep(time.Second) time.Sleep(time.Second)
schedulers.HandleClusterLeaderChange(false) jobServer.HandleClusterLeaderChange(false)
// After running a config change, they should stay off // After running a config change, they should stay off
schedulers.handleConfigChange(nil, nil) jobServer.schedulers.handleConfigChange(nil, nil)
schedulers.Stop() jobServer.StopSchedulers()
for _, element := range schedulers.nextRunTimes { for _, element := range jobServer.schedulers.nextRunTimes {
assert.Nil(t, element) assert.Nil(t, element)
} }
}) })

Просмотреть файл

@@ -4,6 +4,8 @@
package jobs package jobs
import ( import (
"sync"
"github.com/mattermost/mattermost-server/v5/einterfaces" "github.com/mattermost/mattermost-server/v5/einterfaces"
ejobs "github.com/mattermost/mattermost-server/v5/einterfaces/jobs" ejobs "github.com/mattermost/mattermost-server/v5/einterfaces/jobs"
tjobs "github.com/mattermost/mattermost-server/v5/jobs/interfaces" tjobs "github.com/mattermost/mattermost-server/v5/jobs/interfaces"
@@ -16,8 +18,6 @@ type JobServer struct {
ConfigService configservice.ConfigService ConfigService configservice.ConfigService
Store store.Store Store store.Store
metrics einterfaces.MetricsInterface metrics einterfaces.MetricsInterface
Workers *Workers
Schedulers *Schedulers
DataRetentionJob ejobs.DataRetentionJobInterface DataRetentionJob ejobs.DataRetentionJobInterface
MessageExportJob ejobs.MessageExportJobInterface MessageExportJob ejobs.MessageExportJobInterface
@@ -35,6 +35,11 @@ type JobServer struct {
ExportProcess tjobs.ExportProcessInterface ExportProcess tjobs.ExportProcessInterface
ExportDelete tjobs.ExportDeleteInterface ExportDelete tjobs.ExportDeleteInterface
Cloud ejobs.CloudJobInterface 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 { 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() return srv.ConfigService.Config()
} }
func (srv *JobServer) StartWorkers() { func (srv *JobServer) StartWorkers() error {
srv.Workers = srv.Workers.Start() 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() { func (srv *JobServer) StartSchedulers() error {
srv.Schedulers = srv.Schedulers.Start() 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() { func (srv *JobServer) StopWorkers() error {
if srv.Workers != nil { srv.mut.Lock()
srv.Workers.Stop() 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() { func (srv *JobServer) StopSchedulers() error {
if srv.Schedulers != nil { srv.mut.Lock()
srv.Schedulers.Stop() 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)
} }
} }

192
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)
})
}

Просмотреть файл

@@ -4,7 +4,7 @@
package jobs package jobs
import ( import (
"sync" "errors"
"github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/model"
"github.com/mattermost/mattermost-server/v5/services/configservice" "github.com/mattermost/mattermost-server/v5/services/configservice"
@@ -12,7 +12,6 @@ import (
) )
type Workers struct { type Workers struct {
startOnce sync.Once
ConfigService configservice.ConfigService ConfigService configservice.ConfigService
Watcher *Watcher Watcher *Watcher
@@ -34,9 +33,23 @@ type Workers struct {
Cloud model.Worker Cloud model.Worker
listenerId string listenerId string
running bool
}
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
} }
func (srv *JobServer) InitWorkers() *Workers {
workers := &Workers{ workers := &Workers{
ConfigService: srv.ConfigService, ConfigService: srv.ConfigService,
} }
@@ -106,13 +119,15 @@ func (srv *JobServer) InitWorkers() *Workers {
workers.Cloud = cloudInterface.MakeWorker() 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") mlog.Info("Starting workers")
workers.startOnce.Do(func() {
if workers.DataRetention != nil && (*workers.ConfigService.Config().DataRetentionSettings.EnableMessageDeletion || *workers.ConfigService.Config().DataRetentionSettings.EnableFileDeletion) { if workers.DataRetention != nil && (*workers.ConfigService.Config().DataRetentionSettings.EnableMessageDeletion || *workers.ConfigService.Config().DataRetentionSettings.EnableFileDeletion) {
go workers.DataRetention.Run() go workers.DataRetention.Run()
} }
@@ -178,11 +193,9 @@ func (workers *Workers) Start() *Workers {
} }
go workers.Watcher.Start() go workers.Watcher.Start()
})
workers.listenerId = workers.ConfigService.AddConfigListener(workers.handleConfigChange) workers.listenerId = workers.ConfigService.AddConfigListener(workers.handleConfigChange)
workers.running = true
return workers
} }
func (workers *Workers) handleConfigChange(oldConfig *model.Config, newConfig *model.Config) { 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.ConfigService.RemoveConfigListener(workers.listenerId)
workers.Watcher.Stop() workers.Watcher.Stop()
@@ -306,7 +321,7 @@ func (workers *Workers) Stop() *Workers {
workers.Cloud.Stop() workers.Cloud.Stop()
} }
mlog.Info("Stopped workers") workers.running = false
return workers mlog.Info("Stopped workers")
} }