diff --git a/jobs/jobs.go b/jobs/jobs.go index 5e739b3d9d..897fe87c25 100644 --- a/jobs/jobs.go +++ b/jobs/jobs.go @@ -194,7 +194,7 @@ func (srv *JobServer) RequestCancellation(jobId string) *model.AppError { return model.NewAppError("RequestCancellation", "jobs.request_cancellation.status.error", nil, "id="+jobId, http.StatusInternalServerError) } -func (srv *JobServer) CancellationWatcher(ctx context.Context, jobId string, cancelChan chan interface{}) { +func (srv *JobServer) CancellationWatcher(ctx context.Context, jobId string, cancelChan chan struct{}) { for { select { case <-ctx.Done(): @@ -202,11 +202,14 @@ func (srv *JobServer) CancellationWatcher(ctx context.Context, jobId string, can return case <-time.After(CancelWatcherPollingInterval * time.Millisecond): mlog.Debug("CancellationWatcher for Job started polling.", mlog.String("job_id", jobId)) - if jobStatus, err := srv.Store.Job().Get(jobId); err == nil { - if jobStatus.Status == model.JobStatusCancelRequested { - close(cancelChan) - return - } + jobStatus, err := srv.Store.Job().Get(jobId) + if err != nil { + mlog.Warn("Error getting job", mlog.String("job_id", jobId), mlog.Err(err)) + continue + } + if jobStatus.Status == model.JobStatusCancelRequested { + close(cancelChan) + return } } } diff --git a/jobs/migrations/worker.go b/jobs/migrations/worker.go index 1872b68811..001c04a493 100644 --- a/jobs/migrations/worker.go +++ b/jobs/migrations/worker.go @@ -6,6 +6,7 @@ package migrations import ( "context" "net/http" + "sync/atomic" "time" "github.com/mattermost/mattermost-server/v6/jobs" @@ -20,17 +21,18 @@ const ( type Worker struct { name string - stop chan bool + stop chan struct{} stopped chan bool jobs chan model.Job jobServer *jobs.JobServer store store.Store + closed int32 } func MakeWorker(jobServer *jobs.JobServer, store store.Store) model.Worker { worker := Worker{ name: "Migrations", - stop: make(chan bool, 1), + stop: make(chan struct{}), stopped: make(chan bool, 1), jobs: make(chan model.Job), jobServer: jobServer, @@ -41,6 +43,10 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store) model.Worker { } func (worker *Worker) Run() { + // Set to open if closed before. We are not bothered about multiple opens. + if atomic.CompareAndSwapInt32(&worker.closed, 1, 0) { + worker.stop = make(chan struct{}) + } mlog.Debug("Worker started", mlog.String("worker", worker.name)) defer func() { @@ -61,8 +67,12 @@ func (worker *Worker) Run() { } func (worker *Worker) Stop() { + // Set to close, and if already closed before, then return. + if !atomic.CompareAndSwapInt32(&worker.closed, 0, 1) { + return + } mlog.Debug("Worker stopping", mlog.String("worker", worker.name)) - worker.stop <- true + close(worker.stop) <-worker.stopped } @@ -86,7 +96,7 @@ func (worker *Worker) DoJob(job *model.Job) { } cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background()) - cancelWatcherChan := make(chan interface{}, 1) + cancelWatcherChan := make(chan struct{}, 1) go worker.jobServer.CancellationWatcher(cancelCtx, job.Id, cancelWatcherChan) defer cancelCancelWatcher() diff --git a/services/searchengine/bleveengine/indexer/indexing_job.go b/services/searchengine/bleveengine/indexer/indexing_job.go index 510be5ed57..db52e7eaeb 100644 --- a/services/searchengine/bleveengine/indexer/indexing_job.go +++ b/services/searchengine/bleveengine/indexer/indexing_job.go @@ -7,6 +7,7 @@ import ( "context" "net/http" "strconv" + "sync/atomic" "time" "github.com/mattermost/mattermost-server/v6/jobs" @@ -26,11 +27,12 @@ const ( type BleveIndexerWorker struct { name string - stop chan bool + stop chan struct{} stopped chan bool jobs chan model.Job jobServer *jobs.JobServer engine *bleveengine.BleveEngine + closed int32 } func MakeWorker(jobServer *jobs.JobServer, engine *bleveengine.BleveEngine) model.Worker { @@ -39,7 +41,7 @@ func MakeWorker(jobServer *jobs.JobServer, engine *bleveengine.BleveEngine) mode } return &BleveIndexerWorker{ name: "BleveIndexer", - stop: make(chan bool, 1), + stop: make(chan struct{}), stopped: make(chan bool, 1), jobs: make(chan model.Job), jobServer: jobServer, @@ -83,6 +85,10 @@ func (worker *BleveIndexerWorker) IsEnabled(cfg *model.Config) bool { } func (worker *BleveIndexerWorker) Run() { + // Set to open if closed before. We are not bothered about multiple opens. + if atomic.CompareAndSwapInt32(&worker.closed, 1, 0) { + worker.stop = make(chan struct{}) + } mlog.Debug("Worker Started", mlog.String("workername", worker.name)) defer func() { @@ -103,8 +109,12 @@ func (worker *BleveIndexerWorker) Run() { } func (worker *BleveIndexerWorker) Stop() { + // Set to close, and if already closed before, then return. + if !atomic.CompareAndSwapInt32(&worker.closed, 0, 1) { + return + } mlog.Debug("Worker Stopping", mlog.String("workername", worker.name)) - worker.stop <- true + close(worker.stop) <-worker.stopped } @@ -215,7 +225,7 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) { } cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background()) - cancelWatcherChan := make(chan interface{}, 1) + cancelWatcherChan := make(chan struct{}, 1) go worker.jobServer.CancellationWatcher(cancelCtx, job.Id, cancelWatcherChan) defer cancelCancelWatcher()